mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-04 14:25:08 +00:00
Merge remote-tracking branch 'origin/main' into claude/fix-on-demand-edges
# Conflicts: # CHANGELOG.md
This commit is contained in:
Vendored
+40
-5
@@ -4,6 +4,7 @@ Disk Cache
|
||||
Handles persistent disk-based caching with atomic writes and error recovery.
|
||||
"""
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import math
|
||||
import os
|
||||
@@ -31,6 +32,35 @@ except ImportError: # pragma: no cover - exercised on hosts without the wheel
|
||||
# useful, and a half-written file was never useful.
|
||||
_ORPHAN_TEMP_MAX_AGE_SECONDS = 3600
|
||||
|
||||
# Longest key, in UTF-8 bytes, used verbatim as a filename stem. ext4 caps a
|
||||
# name at 255 bytes and set()'s temp file is ".<stem>.json.<8 random>", 15
|
||||
# bytes longer than the stem, so anything near the cap could never be written:
|
||||
# the calendar plugin's key joins every calendar id and passed 300 bytes on a
|
||||
# real install, failing every write with ENAMETOOLONG. Longer keys keep this
|
||||
# many bytes as a readable prefix and end in a hash of the whole key.
|
||||
_MAX_KEY_FILENAME_BYTES = 200
|
||||
_KEY_HASH_CHARS = 16
|
||||
|
||||
|
||||
def _filename_stem(key: str) -> str:
|
||||
"""The filename stem for a key that is already a safe path component.
|
||||
|
||||
Short keys are used as they are, so every file already on disk keeps its
|
||||
name. A long one becomes its first bytes plus a hash of the full key: the
|
||||
prefix keeps the stem recognisable (and keeps the data-type words that
|
||||
cleanup's retention lookup reads from it), the hash keeps two keys that
|
||||
share a long prefix apart. The result is itself short, so a stem read back
|
||||
from a filename -- which is how the web UI names a key it deletes -- maps to
|
||||
the same file.
|
||||
"""
|
||||
encoded = key.encode('utf-8')
|
||||
if len(encoded) <= _MAX_KEY_FILENAME_BYTES:
|
||||
return key
|
||||
digest = hashlib.sha256(encoded).hexdigest()[:_KEY_HASH_CHARS]
|
||||
keep = _MAX_KEY_FILENAME_BYTES - _KEY_HASH_CHARS - 1
|
||||
prefix = encoded[:keep].decode('utf-8', errors='ignore')
|
||||
return f"{prefix}-{digest}"
|
||||
|
||||
|
||||
|
||||
class CacheStrategyProtocol(Protocol):
|
||||
@@ -343,6 +373,8 @@ class DiskCache:
|
||||
derives them), so rejecting anything with a path component turns
|
||||
away only inputs that could never have been written here.
|
||||
|
||||
A key too long to be a filename is shortened by _filename_stem.
|
||||
|
||||
Args:
|
||||
key: Cache key
|
||||
|
||||
@@ -356,7 +388,7 @@ class DiskCache:
|
||||
if safe_key is None:
|
||||
self.logger.warning("Rejected unsafe cache key %r", key)
|
||||
return None
|
||||
return os.path.join(self.cache_dir, f"{safe_key}.json")
|
||||
return os.path.join(self.cache_dir, f"{_filename_stem(safe_key)}.json")
|
||||
|
||||
def get(self, key: str, max_age: Optional[int] = 300) -> Optional[Dict[str, Any]]:
|
||||
"""
|
||||
@@ -561,7 +593,7 @@ class DiskCache:
|
||||
# If direct write also fails, try fallback location
|
||||
self.logger.warning("Direct write failed for key '%s' to %s: %s", key, cache_path, write_error)
|
||||
raise # Re-raise to trigger fallback logic
|
||||
except (IOError, OSError, PermissionError):
|
||||
except (IOError, OSError, PermissionError) as primary_error:
|
||||
# Attempt one-time fallback write to user's home cache directory
|
||||
try:
|
||||
# Try user's home cache directory as fallback
|
||||
@@ -587,11 +619,14 @@ class DiskCache:
|
||||
self.logger.debug("Fallback cache write also failed for key '%s': %s", key, e2)
|
||||
|
||||
# If all write attempts failed, log warning but don't raise exception
|
||||
# Cache is a performance optimization, not critical for operation
|
||||
# Cache is a performance optimization, not critical for operation.
|
||||
# Name the real error: this used to say "permission denied"
|
||||
# whatever happened, which sent a too-long filename off to
|
||||
# be debugged as a directory-ownership problem.
|
||||
self.logger.warning(
|
||||
"Could not write cache for key '%s' to %s (permission denied). "
|
||||
"Could not write cache for key '%s' to %s (%s). "
|
||||
"Cache will be unavailable for this key, but application will continue.",
|
||||
key, cache_path
|
||||
key, cache_path, primary_error.strerror or primary_error
|
||||
)
|
||||
return # Exit gracefully without raising exception
|
||||
|
||||
|
||||
+34
-2
@@ -46,6 +46,32 @@ from src.cache.disk_cache import DateTimeEncoder # noqa: F401 - deliberate re-e
|
||||
# CacheManager.config_manager not built yet (None means "not available").
|
||||
_UNSET: Any = object()
|
||||
|
||||
|
||||
def _outlived(record: Any, max_age: Optional[float], now: float) -> bool:
|
||||
"""Whether a record's own timestamp puts it past max_age.
|
||||
|
||||
The memory tier times an entry from when it was put there, and a record
|
||||
loaded from disk is put there when it is read, not when it was written: a
|
||||
record 290 s old, read after a restart, could be served for another
|
||||
max_age from memory. This is the age check DiskCache.get makes, with the
|
||||
same rule that a stored ttl wins over the caller's max_age. A record that
|
||||
carries no timestamp is left to the memory tier's own clock.
|
||||
"""
|
||||
if not isinstance(record, dict):
|
||||
return False
|
||||
stored_ttl = record.get('ttl')
|
||||
if isinstance(stored_ttl, (int, float)) and not isinstance(stored_ttl, bool) \
|
||||
and stored_ttl >= 0:
|
||||
max_age = stored_ttl
|
||||
stamp = record.get('timestamp')
|
||||
if max_age is None or stamp is None or isinstance(stamp, bool):
|
||||
return False
|
||||
try:
|
||||
return now - float(stamp) > max_age
|
||||
except (TypeError, ValueError):
|
||||
return False
|
||||
|
||||
|
||||
class CacheManager:
|
||||
"""Manages caching of API responses to reduce API calls."""
|
||||
|
||||
@@ -284,7 +310,11 @@ class CacheManager:
|
||||
# 1) Memory cache
|
||||
cached = self._memory_cache_component.get(key, max_age=in_memory_ttl)
|
||||
if cached is not None:
|
||||
return cached
|
||||
if not _outlived(cached, max_age, time.time()):
|
||||
return cached
|
||||
# Too old for this reader. Disk may hold a newer write (from the
|
||||
# other process), and if it does not, the miss is the right answer.
|
||||
self._memory_cache_component.clear(key)
|
||||
|
||||
# 2) Disk cache
|
||||
record = self._disk_cache_component.get(key, max_age=max_age)
|
||||
@@ -318,7 +348,9 @@ class CacheManager:
|
||||
# Check memory cache first (1 minute TTL)
|
||||
cached = self._memory_cache_component.get(key, max_age=60)
|
||||
if cached is not None:
|
||||
return cached
|
||||
if not _outlived(cached, 3600, time.time()):
|
||||
return cached
|
||||
self._memory_cache_component.clear(key)
|
||||
|
||||
# Check disk cache
|
||||
data = self._disk_cache_component.get(key, max_age=3600) # 1 hour for load_cache
|
||||
|
||||
+24
-24
@@ -561,7 +561,7 @@ class ScrollHelper:
|
||||
width = self.display_width
|
||||
strip_width = self.cached_array.shape[1]
|
||||
|
||||
if start_x + width + 1 <= strip_width:
|
||||
if 0 <= start_x and start_x + width + 1 <= strip_width:
|
||||
# Slice the backing array directly. Going via
|
||||
# _get_visible_portion_integer would build two PIL images only for
|
||||
# them to be converted straight back to arrays, which measured 15x
|
||||
@@ -569,9 +569,10 @@ class ScrollHelper:
|
||||
near = self.cached_array[:, start_x:start_x + width]
|
||||
far = self.cached_array[:, start_x + 1:start_x + 1 + width]
|
||||
else:
|
||||
# Close enough to the end that one of the slices wraps; let the
|
||||
# integer path handle that and pay the conversion. Continuous mode
|
||||
# extends the strip before reaching here, so this is the rare case.
|
||||
# One of the slices wraps (close to the end, or a strip narrower
|
||||
# than the panel); let the integer path handle that and pay the
|
||||
# conversion. Continuous mode extends the strip before reaching
|
||||
# here, so this is the rare case.
|
||||
near = np.asarray(
|
||||
self._get_visible_portion_integer(start_x, start_x + width))
|
||||
far = np.asarray(
|
||||
@@ -601,34 +602,33 @@ class ScrollHelper:
|
||||
_size = (self.display_width, self.display_height)
|
||||
img_w = self.cached_array.shape[1]
|
||||
|
||||
if end_x <= img_w:
|
||||
if 0 <= start_x and end_x <= img_w:
|
||||
# Normal case: single contiguous slice (fastest path). tobytes()
|
||||
# on the column-slice view already returns C-order bytes, so
|
||||
# ascontiguousarray() first only added a second full-frame copy.
|
||||
return Image.frombytes(
|
||||
'RGB', _size,
|
||||
self.cached_array[:, start_x:end_x].tobytes())
|
||||
|
||||
# Ensure frame buffer is allocated for all non-simple paths
|
||||
if self._frame_buffer is None or self._frame_buffer.shape != (self.display_height, self.display_width, 3):
|
||||
self._frame_buffer = np.zeros((self.display_height, self.display_width, 3), dtype=np.uint8)
|
||||
|
||||
if img_w == 0:
|
||||
self._frame_buffer[:] = 0
|
||||
else:
|
||||
# Ensure frame buffer is allocated for all non-simple paths
|
||||
if self._frame_buffer is None or self._frame_buffer.shape != (self.display_height, self.display_width, 3):
|
||||
self._frame_buffer = np.zeros((self.display_height, self.display_width, 3), dtype=np.uint8)
|
||||
# The frame runs off the strip, so it carries on from the head:
|
||||
# frame column j is strip column (start_x + j) modulo the strip's
|
||||
# width -- the tail and then the head, and a strip narrower than
|
||||
# the panel repeated across it. Copying the tail and then the rest
|
||||
# of the frame from the head assumed the head was that wide, and
|
||||
# raised at every position for a strip narrower than the panel
|
||||
# (Vegas composes one, with no lead-in, when its content is
|
||||
# narrower than the chain).
|
||||
np.take(self.cached_array, np.arange(start_x, end_x), axis=1,
|
||||
mode='wrap', out=self._frame_buffer)
|
||||
|
||||
width1 = img_w - start_x
|
||||
if width1 > 0:
|
||||
# Wrap-around: tail of image + head of image
|
||||
self._frame_buffer[:, :width1] = self.cached_array[:, start_x:]
|
||||
remaining_width = self.display_width - width1
|
||||
self._frame_buffer[:, width1:] = self.cached_array[:, :remaining_width]
|
||||
else:
|
||||
# Edge case: start_x at or past image end — show from beginning,
|
||||
# clamped to available width (scroll_position should wrap before
|
||||
# reaching this state in normal operation).
|
||||
available = min(self.display_width, img_w)
|
||||
self._frame_buffer[:, :available] = self.cached_array[:, :available]
|
||||
if available < self.display_width:
|
||||
self._frame_buffer[:, available:] = 0
|
||||
|
||||
return Image.frombytes('RGB', _size, self._frame_buffer.tobytes())
|
||||
return Image.frombytes('RGB', _size, self._frame_buffer.tobytes())
|
||||
|
||||
def calculate_dynamic_duration(self) -> int:
|
||||
"""
|
||||
|
||||
@@ -18,7 +18,7 @@ the extra guard only stops a None size raising TypeError.
|
||||
"""
|
||||
|
||||
import logging
|
||||
from datetime import datetime, timezone
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import Any, Dict, Optional, Tuple
|
||||
from zoneinfo import ZoneInfo
|
||||
|
||||
@@ -338,10 +338,46 @@ def format_game_date(config: Optional[Dict[str, Any]], logger, date_text: str,
|
||||
if not raw:
|
||||
return ""
|
||||
fmt = str(scroll_card_option(config, "date_format", "abbrev") or "abbrev")
|
||||
return _format_date_as(fmt, raw, lambda: weekday_for(config, logger, game))
|
||||
return _format_date_as(fmt, raw, lambda: weekday_for(config, logger, game),
|
||||
game=game)
|
||||
|
||||
|
||||
def _format_date_as(fmt: str, raw: str, weekday, months=MONTH_ABBR) -> str:
|
||||
def _printed_weekday(game: Optional[Dict], month: int, day: int) -> str:
|
||||
"""The weekday of the date a card prints as month/day, or '' if unknown.
|
||||
|
||||
The extractor prints "M/D" in the plugin's resolved zone (its own setting,
|
||||
else the global one, else the system zone). The card cannot see that zone:
|
||||
it is handed the plugin's config, whose ``timezone`` ships as "", so
|
||||
card_tzinfo answers UTC and an evening kickoff in the Americas got the
|
||||
next day's weekday ("Sat Oct 2" for a Friday game). Every zone is within
|
||||
a day of UTC, so the printed date is the start's UTC date or a neighbour
|
||||
of it; the one with that month and day is the date on the card.
|
||||
"""
|
||||
if not isinstance(game, dict):
|
||||
return ""
|
||||
raw = game.get("start_time_utc") or game.get("start_time")
|
||||
if not raw:
|
||||
return ""
|
||||
try:
|
||||
start = raw if isinstance(raw, datetime) else datetime.fromisoformat(
|
||||
str(raw).replace("Z", "+00:00"))
|
||||
if start.utcoffset() is None:
|
||||
return "" # naive: no instant to place the date against
|
||||
utc_day = start.astimezone(timezone.utc).date()
|
||||
except (ValueError, TypeError, OverflowError):
|
||||
return ""
|
||||
for offset in (0, -1, 1):
|
||||
try:
|
||||
candidate = utc_day + timedelta(days=offset)
|
||||
except OverflowError:
|
||||
continue
|
||||
if (candidate.month, candidate.day) == (month, day):
|
||||
return WEEKDAY_ABBR[candidate.weekday()]
|
||||
return ""
|
||||
|
||||
|
||||
def _format_date_as(fmt: str, raw: str, weekday, months=MONTH_ABBR,
|
||||
game: Optional[Dict] = None) -> str:
|
||||
"""Render a stripped, non-empty "M/D" *raw* in style *fmt*.
|
||||
|
||||
The body both date formatters share. They differ in which setting names the
|
||||
@@ -349,6 +385,9 @@ def _format_date_as(fmt: str, raw: str, weekday, months=MONTH_ABBR) -> str:
|
||||
``SportsCoreSharedMixin._format_game_date``), so those arrive as arguments:
|
||||
*weekday* is a zero-argument callable, only called for the "weekday" style.
|
||||
*months* lets the mixin keep reading its (overridable) ``_MONTH_ABBR``.
|
||||
With *game*, the "weekday" style names the printed date's own weekday
|
||||
(:func:`_printed_weekday`), and *weekday* is only the fallback for a
|
||||
date its start time cannot place.
|
||||
"""
|
||||
if fmt == "numeric":
|
||||
return raw
|
||||
@@ -364,7 +403,7 @@ def _format_date_as(fmt: str, raw: str, weekday, months=MONTH_ABBR) -> str:
|
||||
if fmt == "day_first":
|
||||
return f"{day} {name}"
|
||||
if fmt == "weekday":
|
||||
day_name = weekday()
|
||||
day_name = _printed_weekday(game, month, day) or weekday()
|
||||
return f"{day_name} {name} {day}" if day_name else f"{name} {day}"
|
||||
return f"{name} {day}"
|
||||
|
||||
|
||||
@@ -360,14 +360,16 @@ class SportsCoreSharedMixin:
|
||||
The formatting is sports_card's. What differs from the card's
|
||||
``format_game_date`` is passed in: the setting (``switch_date_format``,
|
||||
see :meth:`_switch_date_format`) and the weekday, which comes from
|
||||
:meth:`_weekday_for` and so from this plugin's resolved timezone.
|
||||
:meth:`_weekday_for` and so from this plugin's resolved timezone
|
||||
when the game's start cannot place the printed date. The game goes
|
||||
in too, so both formatters name the printed date's own weekday.
|
||||
"""
|
||||
raw = str(date_text or "").strip()
|
||||
if not raw:
|
||||
return raw
|
||||
return _card._format_date_as(self._switch_date_format(), raw,
|
||||
lambda: self._weekday_for(game),
|
||||
self._MONTH_ABBR)
|
||||
self._MONTH_ABBR, game=game)
|
||||
|
||||
def _weekday_for(self, game: Optional[Dict]) -> str:
|
||||
"""Weekday abbreviation from the game's start time, or ''."""
|
||||
|
||||
+92
-42
@@ -14,7 +14,7 @@ import json
|
||||
import time
|
||||
import threading
|
||||
from pathlib import Path
|
||||
from typing import Dict, Any, Optional, List, Callable
|
||||
from typing import Dict, Any, Optional, List, Callable, Tuple
|
||||
from collections import defaultdict
|
||||
import logging
|
||||
import hashlib
|
||||
@@ -52,7 +52,18 @@ class ConfigService:
|
||||
|
||||
# Thread safety
|
||||
self._lock: threading.RLock = threading.RLock()
|
||||
|
||||
# Held across a whole reload -- read, swap, notify -- so one reload's
|
||||
# notifications finish before the next one's start. Subscribers run
|
||||
# under this lock and never under _lock: the display's per-plugin
|
||||
# subscriber can wait seconds for a busy plugin, and get_config(),
|
||||
# subscribe() and unsubscribe() -- called from the render thread --
|
||||
# must not wait behind it.
|
||||
self._notify_lock: threading.RLock = threading.RLock()
|
||||
# (key, callback, thread id) of the callback a notification is running,
|
||||
# so unsubscribe() can wait for that one call; signalled on its return.
|
||||
self._running_callback: Optional[Tuple[str, Callable[..., None], int]] = None
|
||||
self._callback_done = threading.Condition(self._lock)
|
||||
|
||||
# Current configuration
|
||||
self._current_config: Dict[str, Any] = {}
|
||||
self._current_checksum: Optional[str] = None
|
||||
@@ -87,32 +98,33 @@ class ConfigService:
|
||||
True if config changed, False otherwise
|
||||
"""
|
||||
try:
|
||||
new_config = self.config_manager.load_config()
|
||||
new_checksum = self._calculate_checksum(new_config)
|
||||
|
||||
with self._lock:
|
||||
# Check if config actually changed
|
||||
if new_checksum == self._current_checksum:
|
||||
self.logger.debug("Configuration unchanged, skipping reload")
|
||||
return False
|
||||
|
||||
# Store old config for change detection
|
||||
old_config = self._current_config.copy()
|
||||
|
||||
# Update current config
|
||||
self._current_config = new_config
|
||||
self._current_checksum = new_checksum
|
||||
|
||||
# Notify subscribers
|
||||
with self._notify_lock:
|
||||
new_config = self.config_manager.load_config()
|
||||
new_checksum = self._calculate_checksum(new_config)
|
||||
|
||||
with self._lock:
|
||||
# Check if config actually changed
|
||||
if new_checksum == self._current_checksum:
|
||||
self.logger.debug("Configuration unchanged, skipping reload")
|
||||
return False
|
||||
|
||||
# Store old config for change detection
|
||||
old_config = self._current_config.copy()
|
||||
|
||||
# Update current config
|
||||
self._current_config = new_config
|
||||
self._current_checksum = new_checksum
|
||||
|
||||
# Notify subscribers, outside _lock (see _notify_lock)
|
||||
self._notify_subscribers(old_config, new_config)
|
||||
|
||||
|
||||
self.logger.info(
|
||||
"Configuration reloaded (checksum: %s)",
|
||||
new_checksum[:8]
|
||||
)
|
||||
|
||||
|
||||
return True
|
||||
|
||||
|
||||
except ConfigError as e:
|
||||
self.logger.error("Error loading configuration: %s", e, exc_info=True)
|
||||
return False
|
||||
@@ -127,35 +139,64 @@ class ConfigService:
|
||||
Args:
|
||||
old_config: Previous configuration
|
||||
new_config: New configuration
|
||||
|
||||
Called without _lock held. The subscriber lists are copied under it,
|
||||
and each callback is checked against them again just before it runs.
|
||||
"""
|
||||
with self._lock:
|
||||
subscribers = {key: list(callbacks) for key, callbacks in self._subscribers.items()}
|
||||
|
||||
# Notify global subscribers (key: '*')
|
||||
for callback in self._subscribers.get('*', []):
|
||||
try:
|
||||
callback(old_config, new_config)
|
||||
except Exception as e:
|
||||
self.logger.error("Error in global config change callback: %s", e, exc_info=True)
|
||||
|
||||
for callback in subscribers.get('*', []):
|
||||
self._call_subscriber('*', callback, old_config, new_config)
|
||||
|
||||
# Notify plugin-specific subscribers
|
||||
for plugin_id in self._subscribers.keys():
|
||||
for plugin_id, callbacks in subscribers.items():
|
||||
if plugin_id == '*':
|
||||
continue
|
||||
|
||||
|
||||
old_plugin_config = old_config.get(plugin_id, {})
|
||||
new_plugin_config = new_config.get(plugin_id, {})
|
||||
|
||||
|
||||
# Only notify if plugin config actually changed
|
||||
if old_plugin_config != new_plugin_config:
|
||||
for callback in self._subscribers[plugin_id]:
|
||||
try:
|
||||
callback(old_plugin_config, new_plugin_config)
|
||||
except Exception as e:
|
||||
self.logger.error(
|
||||
"Error in config change callback for %s: %s",
|
||||
plugin_id,
|
||||
e,
|
||||
exc_info=True
|
||||
)
|
||||
|
||||
for callback in callbacks:
|
||||
self._call_subscriber(plugin_id, callback,
|
||||
old_plugin_config, new_plugin_config)
|
||||
|
||||
def _call_subscriber(
|
||||
self,
|
||||
key: str,
|
||||
callback: Callable[[Dict[str, Any], Dict[str, Any]], None],
|
||||
old_config: Dict[str, Any],
|
||||
new_config: Dict[str, Any],
|
||||
) -> None:
|
||||
"""Run one callback, unless it was unsubscribed since the snapshot.
|
||||
|
||||
unsubscribe() promises that once it returns the callback is neither
|
||||
running nor will run: the display unloads the plugin straight after.
|
||||
"""
|
||||
with self._lock:
|
||||
if callback not in self._subscribers.get(key, ()):
|
||||
return
|
||||
self._running_callback = (key, callback, threading.get_ident())
|
||||
try:
|
||||
callback(old_config, new_config)
|
||||
except Exception as e:
|
||||
if key == '*':
|
||||
self.logger.error("Error in global config change callback: %s", e, exc_info=True)
|
||||
else:
|
||||
self.logger.error(
|
||||
"Error in config change callback for %s: %s",
|
||||
key,
|
||||
e,
|
||||
exc_info=True
|
||||
)
|
||||
finally:
|
||||
with self._lock:
|
||||
self._running_callback = None
|
||||
self._callback_done.notify_all()
|
||||
|
||||
def _check_file_changes(self) -> bool:
|
||||
"""
|
||||
Check if configuration files have been modified.
|
||||
@@ -276,6 +317,11 @@ class ConfigService:
|
||||
"""
|
||||
Unsubscribe from configuration changes.
|
||||
|
||||
Once this returns the callback is not running and will not be called
|
||||
again. A notification that is running this very callback is waited
|
||||
for (unless the callback is the caller); one running any other
|
||||
callback is not.
|
||||
|
||||
Args:
|
||||
callback: Callback function to remove
|
||||
plugin_id: Optional plugin ID (must match subscription)
|
||||
@@ -285,6 +331,10 @@ class ConfigService:
|
||||
if callback in self._subscribers[key]:
|
||||
self._subscribers[key].remove(callback)
|
||||
self.logger.debug("Unsubscribed from config changes for %s", key)
|
||||
while (self._running_callback is not None
|
||||
and self._running_callback[:2] == (key, callback)
|
||||
and self._running_callback[2] != threading.get_ident()):
|
||||
self._callback_done.wait()
|
||||
|
||||
def shutdown(self) -> None:
|
||||
"""Shutdown the configuration service."""
|
||||
|
||||
@@ -25,11 +25,12 @@ import os
|
||||
import inspect
|
||||
import signal
|
||||
import json
|
||||
import math
|
||||
import threading
|
||||
import types
|
||||
from collections import deque
|
||||
from contextlib import contextmanager
|
||||
from typing import Dict, Any, List, Optional, Callable, Set, Tuple
|
||||
from typing import Dict, Any, FrozenSet, List, Optional, Callable, Set, Tuple
|
||||
from datetime import datetime
|
||||
from concurrent.futures import ThreadPoolExecutor, as_completed # pylint: disable=no-name-in-module
|
||||
import pytz
|
||||
@@ -89,6 +90,19 @@ _MIN_INITIAL_UPDATE_TIMEOUT_SECONDS = 2.0
|
||||
DEFAULT_DYNAMIC_DURATION_CAP = 180.0
|
||||
|
||||
|
||||
def _finite_seconds(value: Any) -> Optional[float]:
|
||||
"""``value`` as seconds when it is a finite number or a numeric string,
|
||||
else None. A bool is not a number here, though it is an int: True would
|
||||
read as a one-second screen."""
|
||||
if isinstance(value, bool):
|
||||
return None
|
||||
try:
|
||||
seconds = float(value)
|
||||
except (TypeError, ValueError, OverflowError):
|
||||
return None
|
||||
return seconds if math.isfinite(seconds) else None
|
||||
|
||||
|
||||
class _PluginReloadJob:
|
||||
"""A ``plugin.reload`` whose slow half runs off the render thread.
|
||||
|
||||
@@ -1348,6 +1362,12 @@ class DisplayController:
|
||||
"until one does", self.EMPTY_ROTATION_PAUSE)
|
||||
self._sleep_with_plugin_updates(self.EMPTY_ROTATION_PAUSE)
|
||||
|
||||
#: Plugins already warned about a display duration that is not a number,
|
||||
#: so a bad setting logs once, not at every one of its screens. A
|
||||
#: frozenset, replaced rather than mutated; class-level default for
|
||||
#: controllers built without __init__ (tests).
|
||||
_duration_warned: FrozenSet[str] = frozenset()
|
||||
|
||||
def _get_display_duration(self, mode_key):
|
||||
"""Seconds to show a mode: the Rotation & Durations page's value for it
|
||||
(display.display_durations), else the plugin's own duration.
|
||||
@@ -1355,6 +1375,17 @@ class DisplayController:
|
||||
The saved value has to win. Every plugin inherits
|
||||
get_display_duration(), so checking the plugin first meant the page's
|
||||
values were never read.
|
||||
|
||||
The plugin's answer is checked here, not trusted. Several plugins
|
||||
return their display_duration setting straight from config.json, so
|
||||
one saved as "20" or null (the raw config editor, a hand edit) came
|
||||
back as a string or None; _resolve_durations compared it with 0, and
|
||||
the TypeError went past every handler in the loop and stopped the
|
||||
display service, which systemd restarted into the same screen. A
|
||||
numeric string counts, as in BasePlugin.get_display_duration; any
|
||||
other value that is not a finite number, or a raise, gets the 30 s a
|
||||
mode without a plugin gets. A number at or below zero is passed on:
|
||||
_resolve_durations has its own rule for that.
|
||||
"""
|
||||
display_durations = self.config.get('display', {}).get('display_durations', {}) or {}
|
||||
override = display_durations.get(mode_key)
|
||||
@@ -1362,8 +1393,22 @@ class DisplayController:
|
||||
return float(override)
|
||||
|
||||
plugin_instance = self.plugin_modes.get(mode_key)
|
||||
if plugin_instance is not None and hasattr(plugin_instance, 'get_display_duration'):
|
||||
return plugin_instance.get_display_duration()
|
||||
if plugin_instance is None or not hasattr(plugin_instance, 'get_display_duration'):
|
||||
return 30
|
||||
try:
|
||||
value = plugin_instance.get_display_duration()
|
||||
except Exception as err: # pylint: disable=broad-except
|
||||
problem = f"get_display_duration() raised {type(err).__name__}: {err}"
|
||||
else:
|
||||
seconds = _finite_seconds(value)
|
||||
if seconds is not None:
|
||||
return seconds
|
||||
problem = f"display duration {value!r} is not a number"
|
||||
plugin_id = getattr(plugin_instance, 'plugin_id', None) or mode_key
|
||||
if plugin_id not in self._duration_warned:
|
||||
self._duration_warned = self._duration_warned | {plugin_id}
|
||||
logger.warning("Plugin %s: %s; showing its modes for 30s (logged once)",
|
||||
plugin_id, problem)
|
||||
return 30
|
||||
|
||||
def _get_global_dynamic_cap(self) -> Optional[float]:
|
||||
|
||||
+6
-1
@@ -396,9 +396,9 @@ class StateSubscription:
|
||||
def _run(self) -> None:
|
||||
backoff = _RECONNECT_MIN_SECONDS
|
||||
while not self._stop.is_set():
|
||||
snapshots = self.snapshots
|
||||
try:
|
||||
self._follow()
|
||||
backoff = _RECONNECT_MIN_SECONDS
|
||||
except ControlError as e:
|
||||
self.last_error = e.reason
|
||||
if e.reason in _SLOW_RETRY_REASONS:
|
||||
@@ -414,6 +414,11 @@ class StateSubscription:
|
||||
sock.close()
|
||||
except OSError:
|
||||
pass
|
||||
if self.snapshots != snapshots:
|
||||
# This connection got as far as the display's state: whatever
|
||||
# ended it (a restart, most often), it was working, so the
|
||||
# next try starts from the shortest wait again.
|
||||
backoff = _RECONNECT_MIN_SECONDS
|
||||
if self._stop.wait(backoff):
|
||||
return
|
||||
backoff = min(backoff * 2, _RECONNECT_MAX_SECONDS)
|
||||
|
||||
@@ -92,7 +92,10 @@ class PluginExecutor:
|
||||
with plugin_scope(plugin_id):
|
||||
result_container['value'] = operation()
|
||||
result_container['completed'] = True
|
||||
except Exception as e:
|
||||
except BaseException as e: # pylint: disable=broad-except
|
||||
# asyncio.CancelledError and SystemExit too: uncaught, one
|
||||
# ended this thread with 'completed' unset, and an operation
|
||||
# that failed at once was reported as timing out.
|
||||
result_container['exception'] = e
|
||||
result_container['completed'] = True
|
||||
|
||||
|
||||
@@ -199,9 +199,22 @@ def contained_plugin_dir(plugin_dir: Path, plugins_dir: Path) -> Optional[str]:
|
||||
name that came out of ``os.scandir()`` on the trusted root carries no
|
||||
taint, which is a real containment guarantee (and one CodeQL's
|
||||
path-injection query can follow), not a string sanitiser.
|
||||
|
||||
The entry looked for is the one ``plugin_dir`` itself names when it sits
|
||||
directly in ``plugins_dir``: for a dev plugin symlinked in under its id,
|
||||
the link's name. Resolving the link first and looking for the target's
|
||||
folder name refused ``plugins/foo -> ~/.ledmatrix-dev-plugins/ledmatrix-foo``
|
||||
(what ``dev_plugin_setup.sh link-github foo <url>`` makes), so the plugin
|
||||
never loaded. Any other path is resolved and matched by its final name,
|
||||
as before.
|
||||
"""
|
||||
plugin_dir_real = os.path.realpath(str(plugin_dir))
|
||||
plugins_dir_real = os.path.realpath(str(plugins_dir))
|
||||
plugin_dir_abs = os.path.abspath(str(plugin_dir))
|
||||
if os.path.realpath(os.path.dirname(plugin_dir_abs)) == plugins_dir_real:
|
||||
matched_name = find_trusted_subdir(plugins_dir_real, os.path.basename(plugin_dir_abs))
|
||||
if matched_name is not None:
|
||||
return os.path.join(plugins_dir_real, matched_name)
|
||||
plugin_dir_real = os.path.realpath(str(plugin_dir))
|
||||
matched_name = find_trusted_subdir(plugins_dir_real, os.path.basename(plugin_dir_real))
|
||||
if matched_name is None:
|
||||
return None
|
||||
@@ -243,6 +256,10 @@ class PluginLoader:
|
||||
self.logger = logger or get_logger(__name__)
|
||||
self._loaded_modules: Dict[str, Any] = {}
|
||||
self._plugin_module_registry: Dict[str, set] = {} # Maps plugin_id to set of module names
|
||||
# plugin_id -> {dotted name: module} for the modules of the plugin's
|
||||
# own packages (``providers.feed``). They keep their names while the
|
||||
# plugin runs and are dropped with it; see _iter_plugin_submodules.
|
||||
self._plugin_submodules: Dict[str, Dict[str, Any]] = {}
|
||||
# Lock to serialize module loading when plugins share module names
|
||||
# (e.g., scroll_display.py, game_renderer.py across sport plugins).
|
||||
# During exec_module, bare-name sub-modules temporarily appear in
|
||||
@@ -449,6 +466,45 @@ class PluginLoader:
|
||||
continue
|
||||
return result
|
||||
|
||||
@staticmethod
|
||||
def _iter_plugin_submodules(
|
||||
plugin_dir: Path, before_keys: set
|
||||
) -> list:
|
||||
"""Return dotted-name modules from plugin_dir added after before_keys.
|
||||
|
||||
The modules of a package the plugin ships (``providers.feed`` from
|
||||
``providers/feed.py``). _iter_plugin_bare_modules skips them, so the
|
||||
bare ``providers`` was namespaced and dropped on unload while
|
||||
``providers.feed`` stayed in sys.modules: a reload after a store update
|
||||
imported a fresh ``providers`` and then got the old ``feed`` back from
|
||||
the cache, running the new manager.py against the old helpers until the
|
||||
display restarted.
|
||||
|
||||
A module counts when its ``__file__`` -- or, for a namespace package,
|
||||
which has none, every ``__path__`` entry -- is inside plugin_dir, so a
|
||||
library the plugin imports (``requests.adapters``) never does.
|
||||
|
||||
Returns a list of (mod_name, module) tuples.
|
||||
"""
|
||||
resolved_dir = plugin_dir.resolve()
|
||||
result = []
|
||||
for key in set(sys.modules.keys()) - before_keys:
|
||||
if "." not in key:
|
||||
continue
|
||||
mod = sys.modules.get(key)
|
||||
if mod is None:
|
||||
continue
|
||||
mod_file = getattr(mod, "__file__", None)
|
||||
locations = [mod_file] if mod_file else list(getattr(mod, "__path__", None) or [])
|
||||
if not locations:
|
||||
continue
|
||||
try:
|
||||
if all(Path(loc).resolve().is_relative_to(resolved_dir) for loc in locations):
|
||||
result.append((key, mod))
|
||||
except (ValueError, TypeError, OSError):
|
||||
continue
|
||||
return result
|
||||
|
||||
def _evict_stale_bare_modules(self, plugin_dir: Path) -> dict:
|
||||
"""Temporarily remove bare-name sys.modules entries from other plugins.
|
||||
|
||||
@@ -527,6 +583,13 @@ class PluginLoader:
|
||||
# Track for cleanup during unload
|
||||
self._plugin_module_registry[plugin_id] = namespaced_names
|
||||
|
||||
# The modules of the plugin's own packages keep their dotted names
|
||||
# while it runs -- as they always have, so the package and its
|
||||
# children stay a matching set in sys.modules -- and are dropped
|
||||
# with the plugin by unregister_plugin_modules().
|
||||
self._plugin_submodules[plugin_id] = dict(
|
||||
self._iter_plugin_submodules(plugin_dir, before_keys))
|
||||
|
||||
if namespaced_names:
|
||||
self.logger.info(
|
||||
"Namespace-isolated %d module(s) for plugin %s",
|
||||
@@ -537,10 +600,16 @@ class PluginLoader:
|
||||
"""Remove namespaced sub-modules and cached module for a plugin from sys.modules.
|
||||
|
||||
Called by PluginManager during unload to clean up all module entries
|
||||
that were created when the plugin was loaded.
|
||||
that were created when the plugin was loaded, including the dotted
|
||||
modules of its packages. A dotted name is dropped only while it still
|
||||
holds this plugin's module: the name is not namespaced, so another
|
||||
plugin may have put its own there since.
|
||||
"""
|
||||
for ns_name in self._plugin_module_registry.pop(plugin_id, set()):
|
||||
sys.modules.pop(ns_name, None)
|
||||
for name, mod in self._plugin_submodules.pop(plugin_id, {}).items():
|
||||
if sys.modules.get(name) is mod:
|
||||
sys.modules.pop(name, None)
|
||||
self._loaded_modules.pop(plugin_id, None)
|
||||
|
||||
def load_module(
|
||||
@@ -646,11 +715,13 @@ class PluginLoader:
|
||||
if evicted_name not in sys.modules:
|
||||
sys.modules[evicted_name] = evicted_mod
|
||||
# Clean up the partially-initialized main module and any
|
||||
# bare-name sub-modules that were added during exec_module
|
||||
# so they don't leak into subsequent plugin loads.
|
||||
# bare-name or package sub-modules that were added during
|
||||
# exec_module so they don't leak into subsequent plugin loads.
|
||||
sys.modules.pop(module_name, None)
|
||||
for key, _ in self._iter_plugin_bare_modules(plugin_dir, before_keys):
|
||||
sys.modules.pop(key, None)
|
||||
for key, _ in self._iter_plugin_submodules(plugin_dir, before_keys):
|
||||
sys.modules.pop(key, None)
|
||||
raise
|
||||
|
||||
self._loaded_modules[plugin_id] = module
|
||||
|
||||
@@ -1163,7 +1163,7 @@ class PluginManager:
|
||||
def _record_update_failure(
|
||||
self,
|
||||
plugin_id: str,
|
||||
exc: Optional[Exception] = None,
|
||||
exc: Optional[BaseException] = None,
|
||||
log: bool = True,
|
||||
count_failure: bool = True,
|
||||
) -> None:
|
||||
@@ -1187,7 +1187,7 @@ class PluginManager:
|
||||
"""
|
||||
failure_time = time.time()
|
||||
if exc is not None:
|
||||
err: Exception = exc
|
||||
err: BaseException = exc
|
||||
error_type = type(exc).__name__
|
||||
else:
|
||||
err = Exception(f"Plugin {plugin_id} execution failed (timeout or executor error)")
|
||||
@@ -1653,7 +1653,7 @@ class PluginManager:
|
||||
finish_guard = threading.Lock()
|
||||
finished = {'done': False}
|
||||
|
||||
def _finish(success: bool, exc: Optional[Exception] = None) -> None:
|
||||
def _finish(success: bool, exc: Optional[BaseException] = None) -> None:
|
||||
with finish_guard:
|
||||
if finished['done']:
|
||||
return
|
||||
@@ -1727,7 +1727,13 @@ class PluginManager:
|
||||
self.resource_monitor.monitor_call(plugin_id, plugin_instance.update)
|
||||
else:
|
||||
plugin_instance.update()
|
||||
except Exception as exc:
|
||||
except BaseException as exc: # pylint: disable=broad-except
|
||||
# BaseException, not just Exception: asyncio.CancelledError
|
||||
# and SystemExit derive from it. Either one skipped _finish,
|
||||
# so the plugin kept its lock and stayed RUNNING for good --
|
||||
# never rescheduled, and every display() skipped as busy.
|
||||
# Re-raised for the executor, which reports it as this
|
||||
# update's failure.
|
||||
_finish(False, exc=exc)
|
||||
raise
|
||||
else:
|
||||
|
||||
@@ -344,12 +344,26 @@ class PluginStoreManager(_RegistryMixin, _InstallMixin, _UpdateMixin):
|
||||
2. Fix permissions via os.chmod() then retry (works for same-owner files)
|
||||
3. Use sudo rm -rf as last resort (works for root-owned __pycache__, etc.)
|
||||
|
||||
A symlink -- a dev plugin linked in by scripts/dev/dev_plugin_setup.sh
|
||||
-- is removed as a link, before any of that: rmtree refuses one, and
|
||||
stage 2 would walk through it and chmod the developer's checkout.
|
||||
|
||||
Args:
|
||||
path: Path to directory to remove
|
||||
|
||||
Returns:
|
||||
True if directory was removed successfully, False otherwise
|
||||
"""
|
||||
if path.is_symlink():
|
||||
# Checked before exists(), which follows the link: a dangling one
|
||||
# would read as already removed and be left behind.
|
||||
try:
|
||||
path.unlink()
|
||||
return True
|
||||
except OSError as e:
|
||||
self.logger.error(f"Could not remove the symlink {path}: {e}")
|
||||
return False
|
||||
|
||||
if not path.exists():
|
||||
return True # Already removed
|
||||
|
||||
|
||||
Reference in New Issue
Block a user