mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-10 09:06:36 +00:00
chore: remove dead code, deprecate unused plugin APIs (over-engineering audit) (#783)
* chore: remove dead code, deprecate unused plugin APIs (over-engineering audit) Whole-tree audit. Every symbol was checked against core, the plugin monorepo and all eight third-party plugins in plugins.json first. - Deprecate (removal 3.10.0) plugin-facing methods nothing calls: LogoDownloader bulk download, ConfigManager backup/secret wrappers, APIHelper extras, BackgroundDataService poll API, PluginManager / PluginStateManager info readers, and a few CacheManager, FontManager, BaseOddsManager, DynamicTeamResolver methods and PluginTestCase. plugin_api_usage.py learns their receiver names; DEPRECATIONS doc regenerated. - Remove core-internal dead code: CacheMetrics, Vegas status/stats plumbing, sync "new cycle" message (followers ignore unknown types), unused operation types, test-only PluginCatalog readers, IPC to_dict and ping, _parse_form_value, CacheStrategyProtocol, ErrorAggregator callbacks, duplicate web response helpers. - Web UI: drop never-mounted json-file-manager.js, the example widget, utils/error_handler.js, four uncalled PluginAPI methods, and 29 escapeHtml shims (call window.LEDEscape directly). Public globals, BaseWidget and widget names unchanged. - Remove six one-off scripts (owner decision) and the unused markupsafe and pytest-mock pins. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(web): calendar picker error text goes in a text node, not innerHTML Same output as the escaped innerHTML it replaces; clears Codacy's XSS-pattern alerts on the line. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -43,6 +43,7 @@ from src.common.espn_dates import (
|
||||
fetch_espn_date_chunks,
|
||||
parse_espn_date_range,
|
||||
)
|
||||
from src.deprecation import deprecated
|
||||
# Configure logging
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@@ -698,6 +699,7 @@ class BackgroundDataService:
|
||||
|
||||
raise last_exception
|
||||
|
||||
@deprecated("3.10.0", "pass callback= to submit_fetch_request()")
|
||||
def get_result(self, request_id: str) -> Optional[FetchResult]:
|
||||
"""
|
||||
Get the result of a fetch request.
|
||||
@@ -714,6 +716,7 @@ class BackgroundDataService:
|
||||
with self._lock:
|
||||
return self.completed_requests.get(request_id)
|
||||
|
||||
@deprecated("3.10.0", "pass callback= to submit_fetch_request()")
|
||||
def is_request_complete(self, request_id: str) -> bool:
|
||||
"""
|
||||
Check if a request has completed.
|
||||
@@ -730,6 +733,7 @@ class BackgroundDataService:
|
||||
with self._lock:
|
||||
return request_id in self.completed_requests
|
||||
|
||||
@deprecated("3.10.0", "pass callback= to submit_fetch_request()")
|
||||
def get_request_status(self, request_id: str) -> Optional[FetchStatus]:
|
||||
"""
|
||||
Get the status of a fetch request.
|
||||
|
||||
@@ -21,6 +21,7 @@ from typing import Dict, Any, Optional, List, cast
|
||||
from src.common.api_helper import DEFAULT_HTTP_HEADERS
|
||||
from src.common.fetch_service import fetch_get, share_connection_pool
|
||||
from src.common.json_body import response_json
|
||||
from src.deprecation import deprecated
|
||||
|
||||
|
||||
|
||||
@@ -277,6 +278,7 @@ class BaseOddsManager:
|
||||
self.logger.warning(f"Unexpected response structure: {json.dumps(data, indent=2)}")
|
||||
return None
|
||||
|
||||
@deprecated("3.10.0", "call get_odds() for each game")
|
||||
def get_odds_for_games(self, games: List[Dict[str, Any]]) -> List[Dict[str, Any]]:
|
||||
"""
|
||||
Fetch odds for multiple games efficiently.
|
||||
@@ -335,6 +337,7 @@ class BaseOddsManager:
|
||||
|
||||
return False
|
||||
|
||||
@deprecated("3.10.0")
|
||||
def format_odds_summary(self, odds_data: Optional[Dict[str, Any]]) -> str:
|
||||
"""
|
||||
Format odds data into a human-readable summary.
|
||||
|
||||
Vendored
-1
@@ -5,6 +5,5 @@ Provides specialized cache components:
|
||||
- MemoryCache: In-memory caching
|
||||
- DiskCache: Persistent disk caching
|
||||
- CacheStrategy: Cache strategy management
|
||||
- CacheMetrics: Performance metrics tracking
|
||||
"""
|
||||
|
||||
|
||||
Vendored
-134
@@ -1,134 +0,0 @@
|
||||
"""
|
||||
Cache Metrics
|
||||
|
||||
Tracks cache performance metrics including hit rates, miss rates, and fetch times.
|
||||
"""
|
||||
|
||||
import threading
|
||||
import time
|
||||
import logging
|
||||
from typing import Dict, Any, Optional
|
||||
|
||||
|
||||
class CacheMetrics:
|
||||
"""Tracks cache performance metrics."""
|
||||
|
||||
def __init__(self, logger: Optional[logging.Logger] = None) -> None:
|
||||
"""
|
||||
Initialize cache metrics tracker.
|
||||
|
||||
Args:
|
||||
logger: Optional logger instance
|
||||
"""
|
||||
self.logger = logger or logging.getLogger(__name__)
|
||||
self._lock = threading.Lock()
|
||||
self._metrics: Dict[str, Any] = {
|
||||
'hits': 0,
|
||||
'misses': 0,
|
||||
'api_calls_saved': 0,
|
||||
'background_hits': 0,
|
||||
'background_misses': 0,
|
||||
'total_fetch_time': 0.0,
|
||||
'fetch_count': 0,
|
||||
# Disk cleanup metrics
|
||||
'last_disk_cleanup': 0.0,
|
||||
'total_files_cleaned': 0,
|
||||
'total_space_freed_mb': 0.0,
|
||||
'last_cleanup_duration_sec': 0.0
|
||||
}
|
||||
|
||||
def record_hit(self, cache_type: str = 'regular') -> None:
|
||||
"""
|
||||
Record a cache hit.
|
||||
|
||||
Args:
|
||||
cache_type: Type of cache hit ('regular' or 'background')
|
||||
"""
|
||||
with self._lock:
|
||||
if cache_type == 'background':
|
||||
self._metrics['background_hits'] += 1
|
||||
else:
|
||||
self._metrics['hits'] += 1
|
||||
|
||||
def record_miss(self, cache_type: str = 'regular') -> None:
|
||||
"""
|
||||
Record a cache miss.
|
||||
|
||||
Args:
|
||||
cache_type: Type of cache miss ('regular' or 'background')
|
||||
"""
|
||||
with self._lock:
|
||||
if cache_type == 'background':
|
||||
self._metrics['background_misses'] += 1
|
||||
else:
|
||||
self._metrics['misses'] += 1
|
||||
self._metrics['api_calls_saved'] += 1
|
||||
|
||||
def record_fetch_time(self, duration: float) -> None:
|
||||
"""
|
||||
Record fetch operation duration.
|
||||
|
||||
Args:
|
||||
duration: Duration in seconds
|
||||
"""
|
||||
with self._lock:
|
||||
self._metrics['total_fetch_time'] += duration
|
||||
self._metrics['fetch_count'] += 1
|
||||
|
||||
def record_disk_cleanup(self, files_cleaned: int, space_freed_mb: float, duration_sec: float) -> None:
|
||||
"""
|
||||
Record disk cleanup operation results.
|
||||
|
||||
Args:
|
||||
files_cleaned: Number of files deleted
|
||||
space_freed_mb: Space freed in megabytes
|
||||
duration_sec: Duration of cleanup operation in seconds
|
||||
"""
|
||||
with self._lock:
|
||||
self._metrics['last_disk_cleanup'] = time.time()
|
||||
self._metrics['total_files_cleaned'] += files_cleaned
|
||||
self._metrics['total_space_freed_mb'] += space_freed_mb
|
||||
self._metrics['last_cleanup_duration_sec'] = duration_sec
|
||||
|
||||
def get_metrics(self) -> Dict[str, Any]:
|
||||
"""
|
||||
Get current cache performance metrics.
|
||||
|
||||
Returns:
|
||||
Dictionary with cache metrics
|
||||
"""
|
||||
with self._lock:
|
||||
total_hits = self._metrics['hits'] + self._metrics['background_hits']
|
||||
total_misses = self._metrics['misses'] + self._metrics['background_misses']
|
||||
total_requests = total_hits + total_misses
|
||||
|
||||
avg_fetch_time = (self._metrics['total_fetch_time'] /
|
||||
self._metrics['fetch_count']) if self._metrics['fetch_count'] > 0 else 0.0
|
||||
|
||||
return {
|
||||
'total_requests': total_requests,
|
||||
'cache_hit_rate': total_hits / total_requests if total_requests > 0 else 0.0,
|
||||
'background_hit_rate': (self._metrics['background_hits'] /
|
||||
(self._metrics['background_hits'] + self._metrics['background_misses'])
|
||||
if (self._metrics['background_hits'] + self._metrics['background_misses']) > 0 else 0.0),
|
||||
'api_calls_saved': self._metrics['api_calls_saved'],
|
||||
'average_fetch_time': avg_fetch_time,
|
||||
'total_fetch_time': self._metrics['total_fetch_time'],
|
||||
'fetch_count': self._metrics['fetch_count'],
|
||||
# Disk cleanup metrics
|
||||
'last_disk_cleanup': self._metrics['last_disk_cleanup'],
|
||||
'total_files_cleaned': self._metrics['total_files_cleaned'],
|
||||
'total_space_freed_mb': self._metrics['total_space_freed_mb'],
|
||||
'last_cleanup_duration_sec': self._metrics['last_cleanup_duration_sec']
|
||||
}
|
||||
|
||||
def log_metrics(self) -> None:
|
||||
"""Log current cache performance metrics."""
|
||||
metrics = self.get_metrics()
|
||||
self.logger.info("Cache Performance - Hit Rate: %.2f%%, Background Hit Rate: %.2f%%, "
|
||||
"API Calls Saved: %d, Avg Fetch Time: %.2fs",
|
||||
metrics['cache_hit_rate'] * 100,
|
||||
metrics['background_hit_rate'] * 100,
|
||||
metrics['api_calls_saved'],
|
||||
metrics['average_fetch_time'])
|
||||
|
||||
Vendored
+7
-34
@@ -4,7 +4,6 @@ Cache Strategy
|
||||
Manages cache strategies (TTLs) for different data types.
|
||||
"""
|
||||
|
||||
import logging
|
||||
from typing import Dict, Any, Optional
|
||||
from datetime import datetime
|
||||
import pytz
|
||||
@@ -13,51 +12,25 @@ import pytz
|
||||
class CacheStrategy:
|
||||
"""Manages cache strategies for different data types."""
|
||||
|
||||
def __init__(self, config_manager: Optional[Any] = None, logger: Optional[logging.Logger] = None) -> None:
|
||||
"""
|
||||
Initialize cache strategy manager.
|
||||
|
||||
Args:
|
||||
config_manager: Optional ConfigManager instance. Kept for callers
|
||||
that pass one; no strategy currently reads it.
|
||||
logger: Optional logger instance
|
||||
"""
|
||||
self.config_manager = config_manager
|
||||
self.logger = logger or logging.getLogger(__name__)
|
||||
|
||||
def get_sport_live_interval(self, sport_key: str) -> int:
|
||||
"""
|
||||
Live-data cache interval, in seconds, for a sport: 60 for every sport.
|
||||
|
||||
This used to read ``live_update_interval`` from a ``<sport>_scoreboard``
|
||||
config section. Those sections belonged to the built-in scoreboards
|
||||
that the plugin system replaced; plugin config is keyed by plugin id
|
||||
(``football-scoreboard``), so the lookup always fell back to 60.
|
||||
|
||||
Args:
|
||||
sport_key: Sport identifier (e.g., 'nba', 'nfl')
|
||||
|
||||
Returns:
|
||||
Live update interval in seconds
|
||||
"""
|
||||
return 60
|
||||
|
||||
def get_cache_strategy(self, data_type: str, sport_key: Optional[str] = None) -> Dict[str, Any]:
|
||||
"""
|
||||
Get cache strategy for different data types.
|
||||
|
||||
Args:
|
||||
data_type: Type of data (e.g., 'live_scores', 'stocks', 'weather_current')
|
||||
sport_key: Optional sport key; for live data it selects the
|
||||
per-sport interval from :meth:`get_sport_live_interval`
|
||||
instead of the generic live default.
|
||||
sport_key: Optional sport key; for live data any sport key
|
||||
selects a 60s interval instead of the generic live default.
|
||||
(That used to be a per-sport ``live_update_interval`` from
|
||||
``<sport>_scoreboard`` config sections, which belonged to the
|
||||
built-in scoreboards the plugin system replaced, so every
|
||||
lookup fell back to 60.)
|
||||
|
||||
Returns:
|
||||
Dictionary with cache strategy (max_age, memory_ttl, etc.)
|
||||
"""
|
||||
live_interval = None
|
||||
if sport_key and data_type in ['sports_live', 'live_scores']:
|
||||
live_interval = self.get_sport_live_interval(sport_key)
|
||||
live_interval = 60
|
||||
|
||||
strategies = {
|
||||
# Ultra time-sensitive data (live scores, current weather)
|
||||
|
||||
Vendored
+6
-20
@@ -15,11 +15,14 @@ import tempfile
|
||||
import logging
|
||||
import threading
|
||||
import zlib
|
||||
from typing import Dict, Any, Optional, Protocol, Tuple
|
||||
from typing import TYPE_CHECKING, Dict, Any, Optional, Tuple
|
||||
from datetime import datetime
|
||||
|
||||
from src.common.path_safety import safe_path_component
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from src.cache.cache_strategy import CacheStrategy
|
||||
|
||||
try: # optional: large speedup on the cache write path, see _dumps below
|
||||
import orjson
|
||||
except ImportError: # pragma: no cover - exercised on hosts without the wheel
|
||||
@@ -62,23 +65,6 @@ def _filename_stem(key: str) -> str:
|
||||
return f"{prefix}-{digest}"
|
||||
|
||||
|
||||
|
||||
class CacheStrategyProtocol(Protocol):
|
||||
"""Protocol for cache strategy objects that categorize cache keys."""
|
||||
|
||||
def get_data_type_from_key(self, key: str) -> str:
|
||||
"""
|
||||
Determine the data type from a cache key.
|
||||
|
||||
Args:
|
||||
key: Cache key
|
||||
|
||||
Returns:
|
||||
Data type string for strategy lookup
|
||||
"""
|
||||
...
|
||||
|
||||
|
||||
class DateTimeEncoder(json.JSONEncoder):
|
||||
"""JSON encoder that handles datetime objects.
|
||||
|
||||
@@ -816,12 +802,12 @@ class DiskCache:
|
||||
# mkstemp's random component.
|
||||
return bool(sep) and len(head) > 1 and bool(suffix)
|
||||
|
||||
def cleanup_expired_files(self, cache_strategy: CacheStrategyProtocol, retention_policies: Dict[str, int]) -> Dict[str, Any]:
|
||||
def cleanup_expired_files(self, cache_strategy: 'CacheStrategy', retention_policies: Dict[str, int]) -> Dict[str, Any]:
|
||||
"""
|
||||
Clean up expired cache files based on retention policies.
|
||||
|
||||
Args:
|
||||
cache_strategy: Object implementing CacheStrategyProtocol for categorizing files
|
||||
cache_strategy: Categorizes files by key (get_data_type_from_key)
|
||||
retention_policies: Dict mapping data types to retention days
|
||||
|
||||
Returns:
|
||||
|
||||
+4
-12
@@ -35,13 +35,13 @@ import tempfile
|
||||
from src.cache.memory_cache import MemoryCache, default_max_size
|
||||
from src.cache.disk_cache import DiskCache
|
||||
from src.cache.cache_strategy import CacheStrategy
|
||||
from src.cache.cache_metrics import CacheMetrics
|
||||
from src.logging_config import get_logger
|
||||
|
||||
# Canonical implementation lives in src.cache.disk_cache; re-exported here
|
||||
# because this module's docstring documents it and external code may import
|
||||
# it from either path.
|
||||
from src.cache.disk_cache import DateTimeEncoder # noqa: F401 - deliberate re-export
|
||||
from src.deprecation import deprecated
|
||||
|
||||
# CacheManager.config_manager not built yet (None means "not available").
|
||||
_UNSET: Any = object()
|
||||
@@ -149,10 +149,7 @@ class CacheManager:
|
||||
max_size=default_max_size(), cleanup_interval=300.0
|
||||
)
|
||||
self._disk_cache_component = DiskCache(cache_dir=self.cache_dir, logger=self.logger)
|
||||
# No config manager: CacheStrategy keeps the parameter for callers but
|
||||
# reads nothing from it, and passing ours would build it eagerly.
|
||||
self._strategy_component = CacheStrategy(logger=self.logger)
|
||||
self._metrics_component = CacheMetrics(logger=self.logger)
|
||||
self._strategy_component = CacheStrategy()
|
||||
|
||||
# Disk cleanup configuration
|
||||
self._disk_cleanup_interval_hours = 24 # Run cleanup every 24 hours
|
||||
@@ -398,6 +395,7 @@ class CacheManager:
|
||||
# caller gets as is.
|
||||
self._disk_cache_component.set(key, data)
|
||||
|
||||
@deprecated("3.10.0", "use get(key, max_age=3600)")
|
||||
def load_cache(self, key: str) -> Optional[Dict[str, Any]]:
|
||||
"""Load data from cache with memory caching."""
|
||||
# Check memory cache first (1 minute TTL)
|
||||
@@ -607,13 +605,6 @@ class CacheManager:
|
||||
duration = time.time() - start_time
|
||||
space_freed_mb = stats['space_freed_bytes'] / (1024 * 1024)
|
||||
|
||||
# Record metrics
|
||||
self._metrics_component.record_disk_cleanup(
|
||||
files_cleaned=stats['files_deleted'],
|
||||
space_freed_mb=space_freed_mb,
|
||||
duration_sec=duration
|
||||
)
|
||||
|
||||
# Log summary
|
||||
if stats['files_deleted'] > 0:
|
||||
self.logger.info(
|
||||
@@ -796,6 +787,7 @@ class CacheManager:
|
||||
data_type = self.get_data_type_from_key(key)
|
||||
return self.get_cached_data_with_strategy(key, data_type)
|
||||
|
||||
@deprecated("3.10.0")
|
||||
def generate_sport_cache_key(self, sport: str, date_str: Optional[str] = None) -> str:
|
||||
"""
|
||||
Centralized cache key generation for sports data.
|
||||
|
||||
@@ -22,6 +22,7 @@ from typing import TYPE_CHECKING, Any, Dict, Mapping, Optional, cast
|
||||
|
||||
import requests
|
||||
from urllib3.util.retry import Retry
|
||||
from src.deprecation import deprecated
|
||||
|
||||
if TYPE_CHECKING:
|
||||
# What Session() puts in .headers; the stubs only promise a MutableMapping.
|
||||
@@ -171,6 +172,7 @@ class APIHelper:
|
||||
self.logger.error(f"Request failed for {url}: {e}")
|
||||
return None
|
||||
|
||||
@deprecated("3.10.0", "use src.common.espn_dates.fetch_espn_scoreboard()")
|
||||
def fetch_espn_scoreboard(self, sport: str, league: str,
|
||||
date: Optional[str] = None,
|
||||
cache_key: Optional[str] = None,
|
||||
@@ -227,6 +229,7 @@ class APIHelper:
|
||||
store_espn_scoreboard_cache(self.cache_manager, shared_key, data)
|
||||
return data
|
||||
|
||||
@deprecated("3.10.0", "call get() with the ESPN URL")
|
||||
def fetch_espn_standings(self, sport: str, league: str,
|
||||
cache_key: Optional[str] = None,
|
||||
cache_ttl: int = 3600) -> Optional[Dict]:
|
||||
@@ -249,6 +252,7 @@ class APIHelper:
|
||||
|
||||
return self.get(url, cache_key=cache_key, cache_ttl=cache_ttl)
|
||||
|
||||
@deprecated("3.10.0", "call get() with the ESPN URL")
|
||||
def fetch_espn_rankings(self, sport: str, league: str,
|
||||
cache_key: Optional[str] = None,
|
||||
cache_ttl: int = 3600) -> Optional[Dict]:
|
||||
@@ -311,6 +315,7 @@ class APIHelper:
|
||||
self.logger.error(f"POST request failed for {url}: {e}")
|
||||
return None
|
||||
|
||||
@deprecated("3.10.0", "use the plugin's cache_manager")
|
||||
def set_cache(self, key: str, data: Any, ttl: int = 3600) -> None:
|
||||
"""
|
||||
Set cache data.
|
||||
@@ -323,6 +328,7 @@ class APIHelper:
|
||||
"""
|
||||
self._set_cache(key, data, ttl)
|
||||
|
||||
@deprecated("3.10.0", "use the plugin's cache_manager")
|
||||
def get_cache(self, key: str) -> Optional[Any]:
|
||||
"""
|
||||
Get cached data.
|
||||
@@ -392,6 +398,7 @@ class APIHelper:
|
||||
self._last_request_monotonic = time.monotonic()
|
||||
self._last_request_time = time.time()
|
||||
|
||||
@deprecated("3.10.0")
|
||||
def set_rate_limit(self, min_interval: float) -> None:
|
||||
"""
|
||||
Set minimum interval between requests.
|
||||
@@ -402,6 +409,7 @@ class APIHelper:
|
||||
self._min_request_interval = min_interval
|
||||
self.logger.debug(f"Rate limit set to {min_interval} seconds")
|
||||
|
||||
@deprecated("3.10.0")
|
||||
def get_request_stats(self) -> Dict[str, Any]:
|
||||
"""
|
||||
Get request statistics.
|
||||
|
||||
@@ -141,7 +141,6 @@ class DisplaySyncManager:
|
||||
self._last_leader_frame_time: float = 0.0
|
||||
self._frame_lock = threading.Lock()
|
||||
self._leader_ip: Optional[str] = None
|
||||
self._on_new_cycle: Optional[Callable[[], None]] = None # called when leader starts new cycle
|
||||
self._on_scroll_image: Optional[Callable[[Image.Image], None]] = None # called with Image when received
|
||||
self._pending_scroll_image: Optional[Image.Image] = None # image received before callback set
|
||||
self._scroll_image_lock = threading.Lock() # guards _on_scroll_image / _pending_scroll_image
|
||||
@@ -412,17 +411,6 @@ class DisplaySyncManager:
|
||||
except Exception as exc:
|
||||
self.logger.debug("Sync: scroll_x send error: %s", exc)
|
||||
|
||||
def send_new_cycle(self) -> None:
|
||||
"""Leader: signal that a new scroll cycle has started so follower rebuilds its image."""
|
||||
if self.role != SyncRole.LEADER:
|
||||
return
|
||||
if self._leader_state != LeaderState.CONNECTED or not self._peer_ip:
|
||||
return
|
||||
try:
|
||||
self._send_sock.sendto(b'{"t":"nc"}', (self._peer_ip, self.port))
|
||||
except Exception as exc:
|
||||
self.logger.debug("Sync: new_cycle send error: %s", exc)
|
||||
|
||||
def send_frame(self, image: Image.Image) -> None:
|
||||
"""Leader: send a rendered frame to the follower as raw RGB bytes.
|
||||
Raw format is orders of magnitude faster than PNG on Pi hardware —
|
||||
@@ -506,21 +494,19 @@ class DisplaySyncManager:
|
||||
self._latest_frame = img
|
||||
self._enter_follower_mode(sender_ip)
|
||||
|
||||
def _enter_follower_mode(self, sender_ip: str) -> bool:
|
||||
def _enter_follower_mode(self, sender_ip: str) -> None:
|
||||
"""Note that the leader at ``sender_ip`` just sent something, and
|
||||
switch from standalone to follower mode if not already following.
|
||||
Returns True if this call made the switch."""
|
||||
switch from standalone to follower mode if not already following."""
|
||||
self._last_leader_frame_time = time.monotonic()
|
||||
self._leader_ip = sender_ip
|
||||
if self._follower_state != FollowerState.STANDALONE:
|
||||
return False
|
||||
return
|
||||
self._follower_state = FollowerState.FOLLOWER
|
||||
self.logger.info(
|
||||
"Sync: leader active at %s — switching to follower mode",
|
||||
sender_ip,
|
||||
)
|
||||
self.write_status_file()
|
||||
return True
|
||||
|
||||
def _follower_recv_loop(self) -> None:
|
||||
while self._running:
|
||||
@@ -570,12 +556,9 @@ class DisplaySyncManager:
|
||||
# frame. Read and validate its fields under a guard —
|
||||
# a UDP payload is attacker-shaped, so a non-object
|
||||
# body makes .get() raise AttributeError and an "sx"
|
||||
# carrying a non-numeric x raises ValueError/TypeError
|
||||
# — but dispatch the callback *outside* it. Running
|
||||
# the callback in here would let a fault in someone
|
||||
# else's code read as a malformed packet and be
|
||||
# logged as one.
|
||||
fire_new_cycle = False
|
||||
# carrying a non-numeric x raises ValueError/TypeError.
|
||||
# Any other "t" is ignored, including the "nc" (new
|
||||
# cycle) that older leaders send and no follower used.
|
||||
try:
|
||||
t = msg.get("t")
|
||||
if t == "hello_ack":
|
||||
@@ -601,18 +584,11 @@ class DisplaySyncManager:
|
||||
# back from. Treat it as malformed.
|
||||
raise ValueError(f"non-finite scroll x: {msg['x']!r}")
|
||||
self._latest_scroll_x = scroll_x
|
||||
if self._enter_follower_mode(sender_ip):
|
||||
fire_new_cycle = True # build initial scroll image
|
||||
elif t == "nc":
|
||||
# Leader started a new scroll cycle — rebuild local image
|
||||
fire_new_cycle = True
|
||||
self._enter_follower_mode(sender_ip)
|
||||
except (KeyError, AttributeError, TypeError, ValueError) as exc:
|
||||
self.logger.debug("Sync: malformed control message: %s", exc)
|
||||
continue
|
||||
|
||||
if fire_new_cycle and self._on_new_cycle:
|
||||
self._on_new_cycle()
|
||||
|
||||
except socket.timeout:
|
||||
continue
|
||||
except Exception as exc:
|
||||
@@ -679,15 +655,6 @@ class DisplaySyncManager:
|
||||
"""Follower: return the most recently received Vegas scroll position, or None."""
|
||||
return self._latest_scroll_x
|
||||
|
||||
def set_on_new_cycle(self, callback: Callable[[], None]) -> None:
|
||||
"""Follower: register a callback fired when the leader starts a new scroll cycle.
|
||||
|
||||
Nothing in core registers one: display_controller follows the leader
|
||||
through set_on_scroll_image() and the scroll position instead of
|
||||
rebuilding locally. The hook stays for callers that want the signal.
|
||||
"""
|
||||
self._on_new_cycle = callback
|
||||
|
||||
def get_latest_frame(self) -> Optional[Image.Image]:
|
||||
"""Follower: return the most recently received pixel frame (non-Vegas fallback)."""
|
||||
with self._frame_lock:
|
||||
|
||||
@@ -45,6 +45,7 @@ from src.common.permission_utils import (
|
||||
ensure_shared_group_ownership,
|
||||
get_config_dir_mode
|
||||
)
|
||||
from src.deprecation import deprecated
|
||||
|
||||
|
||||
def _private_copy(config: Dict[str, Any]) -> Dict[str, Any]:
|
||||
@@ -174,6 +175,7 @@ class ConfigManager:
|
||||
|
||||
return result
|
||||
|
||||
@deprecated("3.10.0", "backups are handled by src.backup_manager")
|
||||
def rollback_config(self, backup_version: Optional[str] = None) -> bool:
|
||||
"""
|
||||
Rollback configuration to a previous backup.
|
||||
@@ -198,6 +200,7 @@ class ConfigManager:
|
||||
|
||||
return success
|
||||
|
||||
@deprecated("3.10.0", "backups are handled by src.backup_manager")
|
||||
def list_backups(self) -> List[BackupInfo]:
|
||||
"""
|
||||
List all available configuration backups.
|
||||
@@ -208,6 +211,7 @@ class ConfigManager:
|
||||
atomic_mgr = self._get_atomic_manager()
|
||||
return atomic_mgr.list_backups()
|
||||
|
||||
@deprecated("3.10.0")
|
||||
def validate_config_file(self, config_path: Optional[str] = None) -> ValidationResult:
|
||||
"""
|
||||
Validate a configuration file.
|
||||
@@ -404,6 +408,7 @@ class ConfigManager:
|
||||
self.logger.error(error_msg, exc_info=True)
|
||||
raise ConfigError(error_msg, config_path=self.config_path) from e
|
||||
|
||||
@deprecated("3.10.0", "secrets are merged into each plugin's config; read them with config.get()")
|
||||
def get_secret(self, key: str) -> Optional[Any]:
|
||||
"""Get a secret value by key."""
|
||||
try:
|
||||
@@ -757,6 +762,7 @@ class ConfigManager:
|
||||
self.logger.error(error_msg, exc_info=True)
|
||||
raise ConfigError(error_msg, config_path=self.config_path, field=plugin_id) from e
|
||||
|
||||
@deprecated("3.10.0")
|
||||
def cleanup_orphaned_plugin_configs(self, valid_plugin_ids: List[str]) -> List[str]:
|
||||
"""
|
||||
Remove configuration sections for plugins that are no longer installed.
|
||||
@@ -810,6 +816,7 @@ class ConfigManager:
|
||||
self.logger.error(f"Error cleaning up orphaned plugin configs: {e}")
|
||||
return removed
|
||||
|
||||
@deprecated("3.10.0")
|
||||
def validate_all_plugin_configs(self, plugin_schema_manager=None) -> Dict[str, Dict[str, Any]]:
|
||||
"""
|
||||
Validate all plugin configurations against their schemas.
|
||||
|
||||
@@ -24,6 +24,7 @@ from typing import Any, Dict, List
|
||||
|
||||
from src.common.api_helper import DEFAULT_HTTP_HEADERS
|
||||
from src.common.json_body import response_json
|
||||
from src.deprecation import deprecated
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@@ -201,6 +202,7 @@ class DynamicTeamResolver:
|
||||
DynamicTeamResolver._failure_timestamp = current_time
|
||||
return {}
|
||||
|
||||
@deprecated("3.10.0", "use resolve_teams()")
|
||||
def get_available_dynamic_teams(self) -> List[str]:
|
||||
"""
|
||||
Get list of available dynamic team names.
|
||||
@@ -210,6 +212,7 @@ class DynamicTeamResolver:
|
||||
"""
|
||||
return list(self.DYNAMIC_PATTERNS.keys())
|
||||
|
||||
@deprecated("3.10.0", "use resolve_teams()")
|
||||
def is_dynamic_team(self, team_name: str) -> bool:
|
||||
"""
|
||||
Check if a team name is a dynamic team.
|
||||
|
||||
+2
-40
@@ -123,8 +123,7 @@ class ErrorAggregator:
|
||||
self._error_counts: Dict[str, int] = defaultdict(int)
|
||||
self._plugin_error_counts: Dict[str, Dict[str, int]] = defaultdict(lambda: defaultdict(int))
|
||||
self._patterns: Dict[str, ErrorPattern] = {}
|
||||
self._pattern_callbacks: List[Callable[[ErrorPattern], None]] = []
|
||||
self._lock = threading.RLock() # RLock: build_snapshot and pattern callbacks re-enter
|
||||
self._lock = threading.RLock() # RLock: build_snapshot re-enters
|
||||
|
||||
# Track session start for relative timing
|
||||
self._session_start = datetime.now()
|
||||
@@ -238,13 +237,6 @@ class ErrorAggregator:
|
||||
f"{count} times in last {self.pattern_window}. "
|
||||
f"Affected plugins: {set(affected_plugins) or 'unknown'}"
|
||||
)
|
||||
|
||||
# Notify callbacks
|
||||
for callback in self._pattern_callbacks:
|
||||
try:
|
||||
callback(pattern)
|
||||
except Exception as e:
|
||||
self.logger.error(f"Pattern callback failed: {e}")
|
||||
else:
|
||||
# Update existing pattern
|
||||
self._patterns[pattern_key].count = count
|
||||
@@ -253,15 +245,6 @@ class ErrorAggregator:
|
||||
known = self._patterns[pattern_key].affected_plugins
|
||||
known.extend(p for p in affected_plugins if p not in known)
|
||||
|
||||
def on_pattern_detected(self, callback: Callable[[ErrorPattern], None]) -> None:
|
||||
"""
|
||||
Register a callback to be called when a new error pattern is detected.
|
||||
|
||||
Args:
|
||||
callback: Function that takes an ErrorPattern as argument
|
||||
"""
|
||||
self._pattern_callbacks.append(callback)
|
||||
|
||||
def get_error_summary(self) -> Dict[str, Any]:
|
||||
"""
|
||||
Get summary of all errors for reporting.
|
||||
@@ -325,27 +308,6 @@ class ErrorAggregator:
|
||||
"last_error": recent_plugin_errors[-1].to_dict() if recent_plugin_errors else None
|
||||
}
|
||||
|
||||
def clear_old_records(self, max_age_hours: int = 24) -> int:
|
||||
"""
|
||||
Clear records older than specified age.
|
||||
|
||||
Args:
|
||||
max_age_hours: Maximum age in hours
|
||||
|
||||
Returns:
|
||||
Number of records cleared
|
||||
"""
|
||||
with self._lock:
|
||||
cutoff = datetime.now() - timedelta(hours=max_age_hours)
|
||||
original_count = len(self._records)
|
||||
self._records = [r for r in self._records if r.timestamp > cutoff]
|
||||
cleared = original_count - len(self._records)
|
||||
|
||||
if cleared > 0:
|
||||
self.logger.info(f"Cleared {cleared} old error records")
|
||||
|
||||
return cleared
|
||||
|
||||
@property
|
||||
def version(self) -> int:
|
||||
"""Changes whenever the recorded errors do (see ErrorSnapshotPublisher)."""
|
||||
@@ -354,7 +316,7 @@ class ErrorAggregator:
|
||||
def clear_before(self, cutoff: datetime) -> int:
|
||||
"""Forget every error recorded at or before ``cutoff``.
|
||||
|
||||
Unlike clear_old_records, this also resets what the summary reports:
|
||||
This also resets what the summary reports:
|
||||
the per-type and per-plugin counts are rebuilt from the records that
|
||||
remain, and detected patterns that began before the cutoff are dropped
|
||||
(one that is still happening is detected again on its next
|
||||
|
||||
@@ -41,6 +41,7 @@ from PIL import ImageFont
|
||||
from src.common.bdf_font import load_bdf_face, read_bdf_native_size
|
||||
from src.common.font_layout import load_truetype, resolve_asset_path
|
||||
from typing import Dict, Tuple, Optional, Union, Any
|
||||
from src.deprecation import deprecated
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@@ -533,6 +534,7 @@ class FontManager:
|
||||
"""
|
||||
return load_bdf_face(font_path, size_px)[0]
|
||||
|
||||
@deprecated("3.10.0", "use src.common.bdf_font.read_bdf_native_size()")
|
||||
def get_native_bdf_size(self, family: str) -> Optional[int]:
|
||||
"""The one true pixel size of a BDF family in the catalog, or None
|
||||
for scalable (TTF) families / unknown families."""
|
||||
@@ -553,6 +555,7 @@ class FontManager:
|
||||
|
||||
# ==================== Font Measurement ====================
|
||||
|
||||
@deprecated("3.10.0", "use src.adaptive_layout.measure_ink()")
|
||||
def measure_text(self, text: str, font: Union[ImageFont.FreeTypeFont, freetype.Face]) -> Tuple[int, int, int]:
|
||||
"""
|
||||
Measure text dimensions and baseline.
|
||||
|
||||
@@ -288,11 +288,6 @@ def errors_clear(request_id: str, cutoff: float, *,
|
||||
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)
|
||||
|
||||
|
||||
def hello(client: str = 'web', *, timeout: float = DEFAULT_TIMEOUT_SECONDS,
|
||||
paths: Optional[Sequence[str]] = None) -> Dict[str, Any]:
|
||||
"""Version negotiation: the result's ``version`` is the one both sides speak."""
|
||||
|
||||
@@ -399,9 +399,6 @@ class HelloArgs:
|
||||
versions: Tuple[int, ...] = (PROTOCOL_VERSION,)
|
||||
client: str = ''
|
||||
|
||||
def to_dict(self) -> Dict[str, Any]:
|
||||
return {'versions': list(self.versions), 'client': self.client}
|
||||
|
||||
@classmethod
|
||||
def from_dict(cls, args: Mapping[str, Any]) -> 'HelloArgs':
|
||||
versions = args.get('versions', [PROTOCOL_VERSION])
|
||||
@@ -426,10 +423,6 @@ class OnDemandStartArgs:
|
||||
duration: Optional[float] = None
|
||||
pinned: bool = False
|
||||
|
||||
def to_dict(self) -> Dict[str, Any]:
|
||||
return {'plugin_id': self.plugin_id, 'mode': self.mode,
|
||||
'duration': self.duration, 'pinned': self.pinned}
|
||||
|
||||
@classmethod
|
||||
def from_dict(cls, args: Mapping[str, Any]) -> 'OnDemandStartArgs':
|
||||
plugin_id = _optional_name(args, 'plugin_id')
|
||||
@@ -449,9 +442,6 @@ class OnDemandStartArgs:
|
||||
class OnDemandStopArgs:
|
||||
"""``on_demand.stop``: end the on-demand session and resume rotation."""
|
||||
|
||||
def to_dict(self) -> Dict[str, Any]:
|
||||
return {}
|
||||
|
||||
@classmethod
|
||||
def from_dict(cls, args: Mapping[str, Any]) -> 'OnDemandStopArgs':
|
||||
return cls()
|
||||
@@ -461,9 +451,6 @@ class OnDemandStopArgs:
|
||||
class NoArgs:
|
||||
"""``ping`` and ``on_demand.status`` take no arguments (extra ones are ignored)."""
|
||||
|
||||
def to_dict(self) -> Dict[str, Any]:
|
||||
return {}
|
||||
|
||||
@classmethod
|
||||
def from_dict(cls, args: Mapping[str, Any]) -> 'NoArgs':
|
||||
return cls()
|
||||
@@ -480,9 +467,6 @@ class BrightnessSetArgs:
|
||||
"""
|
||||
brightness: int
|
||||
|
||||
def to_dict(self) -> Dict[str, Any]:
|
||||
return {'brightness': self.brightness}
|
||||
|
||||
@classmethod
|
||||
def from_dict(cls, args: Mapping[str, Any]) -> 'BrightnessSetArgs':
|
||||
value = args.get('brightness')
|
||||
@@ -503,9 +487,6 @@ class PluginReloadArgs:
|
||||
"""
|
||||
plugin_id: str
|
||||
|
||||
def to_dict(self) -> Dict[str, Any]:
|
||||
return {'plugin_id': self.plugin_id}
|
||||
|
||||
@classmethod
|
||||
def from_dict(cls, args: Mapping[str, Any]) -> 'PluginReloadArgs':
|
||||
plugin_id = _optional_name(args, 'plugin_id')
|
||||
@@ -546,9 +527,6 @@ class StateGetArgs:
|
||||
since: Optional[int] = None
|
||||
epoch: Optional[str] = None
|
||||
|
||||
def to_dict(self) -> Dict[str, Any]:
|
||||
return {'since': self.since, 'epoch': self.epoch}
|
||||
|
||||
@classmethod
|
||||
def from_dict(cls, args: Mapping[str, Any]) -> 'StateGetArgs':
|
||||
return cls(since=_optional_version(args, 'since'), epoch=_optional_epoch(args))
|
||||
@@ -565,9 +543,6 @@ class StateSubscribeArgs:
|
||||
and a ``tick`` at least every :data:`SUBSCRIBE_KEEPALIVE_SECONDS`.
|
||||
"""
|
||||
|
||||
def to_dict(self) -> Dict[str, Any]:
|
||||
return {}
|
||||
|
||||
@classmethod
|
||||
def from_dict(cls, args: Mapping[str, Any]) -> 'StateSubscribeArgs':
|
||||
return cls()
|
||||
@@ -582,9 +557,6 @@ class ErrorsClearArgs:
|
||||
"""
|
||||
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')
|
||||
|
||||
@@ -29,6 +29,7 @@ from src.common.permission_utils import (
|
||||
get_assets_dir_mode,
|
||||
get_assets_file_mode
|
||||
)
|
||||
from src.deprecation import deprecated
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@@ -471,6 +472,7 @@ class LogoDownloader:
|
||||
logger.info(f"Using dynamic ESPN endpoint for custom soccer league: {league}")
|
||||
return api_url
|
||||
|
||||
@deprecated("3.10.0", "download logos one at a time with download_missing_logo()")
|
||||
def fetch_teams_data(self, league: str) -> Optional[Dict]:
|
||||
"""Fetch team data from ESPN API for a specific league."""
|
||||
api_url = self._resolve_api_url(league)
|
||||
@@ -518,6 +520,7 @@ class LogoDownloader:
|
||||
logger.error(f"Error parsing JSON response for {team_id} in {league}: {e}")
|
||||
return None
|
||||
|
||||
@deprecated("3.10.0", "download logos one at a time with download_missing_logo()")
|
||||
def extract_teams_from_data(self, data: Dict, league: str) -> List[Dict[str, str]]:
|
||||
"""Extract team information from ESPN API response."""
|
||||
teams = []
|
||||
@@ -621,6 +624,7 @@ class LogoDownloader:
|
||||
# Default to FBS for unknown conferences
|
||||
return 'FBS'
|
||||
|
||||
@deprecated("3.10.0", "download logos one at a time with download_missing_logo()")
|
||||
def download_missing_logos_for_league(self, league: str, force_download: bool = False) -> Tuple[int, int]:
|
||||
"""Download missing logos for a specific league."""
|
||||
logger.info(f"Starting logo download for league: {league}")
|
||||
@@ -675,6 +679,7 @@ class LogoDownloader:
|
||||
logger.info(f"Logo download complete for {league}: {downloaded_count} downloaded, {failed_count} failed")
|
||||
return downloaded_count, failed_count
|
||||
|
||||
@deprecated("3.10.0", "download logos one at a time with download_missing_logo()")
|
||||
def download_all_ncaa_football_logos(self, include_fcs: bool = True, force_download: bool = False) -> Tuple[int, int]:
|
||||
"""Download all NCAA football team logos including FCS teams."""
|
||||
logger.info(f"Starting comprehensive NCAA football logo download (FCS: {include_fcs})")
|
||||
@@ -763,6 +768,7 @@ class LogoDownloader:
|
||||
time.sleep(0.1) # Small delay
|
||||
return success
|
||||
|
||||
@deprecated("3.10.0", "download logos one at a time with download_missing_logo()")
|
||||
def download_all_missing_logos(self, leagues: List[str] | None = None, force_download: bool = False) -> Dict[str, Tuple[int, int]]:
|
||||
"""Download missing logos for all specified leagues."""
|
||||
if leagues is None:
|
||||
@@ -858,6 +864,7 @@ class LogoDownloader:
|
||||
logger.error(f"Failed to create placeholder logo for {team_abbreviation}: {e}")
|
||||
return False
|
||||
|
||||
@deprecated("3.10.0")
|
||||
def convert_image_to_rgba(self, filepath: Path) -> bool:
|
||||
"""Convert an image file to RGBA format to avoid PIL warnings."""
|
||||
try:
|
||||
@@ -875,6 +882,7 @@ class LogoDownloader:
|
||||
logger.error(f"Failed to convert {filepath.name} to RGBA: {e}")
|
||||
return False
|
||||
|
||||
@deprecated("3.10.0")
|
||||
def convert_all_logos_to_rgba(self, league: str) -> Tuple[int, int]:
|
||||
"""Convert all logos in a league directory to RGBA format."""
|
||||
logo_dir = Path(self.get_logo_directory(league))
|
||||
|
||||
@@ -51,44 +51,36 @@ class OperationHistory:
|
||||
def __init__(
|
||||
self,
|
||||
history_file: Optional[str] = None,
|
||||
max_records: int = 1000,
|
||||
lazy_load: bool = False
|
||||
max_records: int = 1000
|
||||
):
|
||||
"""
|
||||
Initialize operation history.
|
||||
|
||||
Initialize operation history. The history file is read on first use,
|
||||
not here, so constructing this costs the web app's startup nothing.
|
||||
|
||||
Args:
|
||||
history_file: Path to file for persisting history
|
||||
max_records: Maximum number of records to keep
|
||||
lazy_load: If True, defer loading history file until first access
|
||||
"""
|
||||
self.logger = get_logger(__name__)
|
||||
self.history_file = Path(history_file) if history_file else None
|
||||
self.max_records = max_records
|
||||
self._lazy_load = lazy_load
|
||||
self._history_loaded = False
|
||||
|
||||
|
||||
# In-memory history
|
||||
self._history: List[OperationRecord] = []
|
||||
self._lock = threading.RLock()
|
||||
|
||||
# Load history from file if it exists (unless lazy loading)
|
||||
if not self._lazy_load and self.history_file and self.history_file.exists():
|
||||
self._load_history()
|
||||
self._history_loaded = True
|
||||
|
||||
|
||||
def _ensure_loaded(self) -> None:
|
||||
"""Ensure history is loaded (for lazy loading)."""
|
||||
"""Load the history file on first use."""
|
||||
if not self._history_loaded and self.history_file and self.history_file.exists():
|
||||
self._load_history()
|
||||
self._history_loaded = True
|
||||
|
||||
|
||||
def record_operation(
|
||||
self,
|
||||
operation_type: str,
|
||||
plugin_id: Optional[str] = None,
|
||||
status: str = "completed",
|
||||
user: Optional[str] = None,
|
||||
details: Optional[Dict[str, Any]] = None,
|
||||
error: Optional[str] = None,
|
||||
operation_id: Optional[str] = None
|
||||
@@ -100,7 +92,6 @@ class OperationHistory:
|
||||
operation_type: Type of operation (install, update, uninstall, etc.)
|
||||
plugin_id: Plugin identifier
|
||||
status: Operation status
|
||||
user: User who performed operation
|
||||
details: Optional operation details
|
||||
error: Optional error message
|
||||
operation_id: Optional operation ID
|
||||
@@ -118,7 +109,6 @@ class OperationHistory:
|
||||
plugin_id=plugin_id,
|
||||
timestamp=datetime.now(),
|
||||
status=status,
|
||||
user=user,
|
||||
details=details,
|
||||
error=error
|
||||
)
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
Plugin operation queue manager.
|
||||
|
||||
Serializes plugin operations to prevent conflicts and provides
|
||||
status tracking and cancellation support.
|
||||
status tracking.
|
||||
"""
|
||||
|
||||
import threading
|
||||
@@ -25,8 +25,8 @@ class PluginOperationQueue:
|
||||
- Serialized execution (one operation at a time)
|
||||
- Prevents concurrent operations on same plugin
|
||||
- Operation status tracking
|
||||
- Operation cancellation
|
||||
- In-memory history of finished operations
|
||||
- A bounded in-memory history of finished operations, which also caps
|
||||
how many finished operations get_operation_status() remembers
|
||||
|
||||
The history is not persisted. The web UI's operation history comes from
|
||||
OperationHistory (operation_history.py), which has its own file; a copy
|
||||
@@ -133,56 +133,6 @@ class PluginOperationQueue:
|
||||
with self._lock:
|
||||
return self._operations.get(operation_id)
|
||||
|
||||
def cancel_operation(self, operation_id: str) -> bool:
|
||||
"""
|
||||
Cancel a pending operation.
|
||||
|
||||
Args:
|
||||
operation_id: Operation identifier
|
||||
|
||||
Returns:
|
||||
True if operation was cancelled, False if not found or already running
|
||||
"""
|
||||
with self._lock:
|
||||
operation = self._operations.get(operation_id)
|
||||
if not operation:
|
||||
return False
|
||||
|
||||
if operation.status == OperationStatus.RUNNING:
|
||||
self.logger.warning(
|
||||
f"Cannot cancel running operation {operation_id}"
|
||||
)
|
||||
return False
|
||||
|
||||
if operation.status == OperationStatus.PENDING:
|
||||
operation.status = OperationStatus.CANCELLED
|
||||
operation.completed_at = datetime.now()
|
||||
operation.message = "Operation cancelled by user"
|
||||
self._add_to_history(operation)
|
||||
self.logger.info(f"Cancelled operation {operation_id}")
|
||||
return True
|
||||
|
||||
return False
|
||||
|
||||
def get_operation_history(self, limit: int = 50) -> List[PluginOperation]:
|
||||
"""
|
||||
Get operation history.
|
||||
|
||||
Args:
|
||||
limit: Maximum number of operations to return
|
||||
|
||||
Returns:
|
||||
List of operations, sorted by creation time (newest first)
|
||||
"""
|
||||
with self._lock:
|
||||
# Sort by creation time (newest first)
|
||||
history = sorted(
|
||||
self._operation_history,
|
||||
key=lambda op: op.created_at,
|
||||
reverse=True
|
||||
)
|
||||
return history[:limit]
|
||||
|
||||
def _start_worker(self) -> None:
|
||||
"""Start the worker thread that processes operations."""
|
||||
if self._worker_thread and self._worker_thread.is_alive():
|
||||
@@ -207,11 +157,6 @@ class PluginOperationQueue:
|
||||
except queue.Empty:
|
||||
continue
|
||||
|
||||
# Check if operation was cancelled
|
||||
if operation.status == OperationStatus.CANCELLED:
|
||||
self._operation_queue.task_done()
|
||||
continue
|
||||
|
||||
# Execute operation
|
||||
self._execute_operation(operation)
|
||||
|
||||
|
||||
@@ -15,11 +15,7 @@ import uuid
|
||||
class OperationType(Enum):
|
||||
"""Types of plugin operations."""
|
||||
INSTALL = "install"
|
||||
UPDATE = "update"
|
||||
UNINSTALL = "uninstall"
|
||||
ENABLE = "enable"
|
||||
DISABLE = "disable"
|
||||
CONFIGURE = "configure"
|
||||
|
||||
|
||||
class OperationStatus(Enum):
|
||||
@@ -28,7 +24,6 @@ class OperationStatus(Enum):
|
||||
RUNNING = "running"
|
||||
COMPLETED = "completed"
|
||||
FAILED = "failed"
|
||||
CANCELLED = "cancelled"
|
||||
|
||||
|
||||
@dataclass
|
||||
@@ -70,29 +65,3 @@ class PluginOperation:
|
||||
'started_at': self.started_at.isoformat() if self.started_at else None,
|
||||
'completed_at': self.completed_at.isoformat() if self.completed_at else None,
|
||||
}
|
||||
|
||||
@classmethod
|
||||
def from_dict(cls, data: Dict[str, Any]) -> 'PluginOperation':
|
||||
"""Create operation from dictionary."""
|
||||
op = cls(
|
||||
operation_type=OperationType(data['operation_type']),
|
||||
plugin_id=data['plugin_id'],
|
||||
operation_id=data.get('operation_id', str(uuid.uuid4())),
|
||||
parameters=data.get('parameters', {}),
|
||||
status=OperationStatus(data.get('status', 'pending')),
|
||||
progress=data.get('progress', 0.0),
|
||||
message=data.get('message', ''),
|
||||
error=data.get('error'),
|
||||
result=data.get('result'),
|
||||
)
|
||||
|
||||
# Parse datetime fields
|
||||
if data.get('created_at'):
|
||||
op.created_at = datetime.fromisoformat(data['created_at'])
|
||||
if data.get('started_at'):
|
||||
op.started_at = datetime.fromisoformat(data['started_at'])
|
||||
if data.get('completed_at'):
|
||||
op.completed_at = datetime.fromisoformat(data['completed_at'])
|
||||
|
||||
return op
|
||||
|
||||
|
||||
@@ -3,9 +3,10 @@ Plugin catalog: what the web process knows about installed plugins.
|
||||
|
||||
The web interface and the display run as two processes. Only the display
|
||||
imports plugin code and runs it; the web process reads plugins as files --
|
||||
manifest, config schema, the plugin's section of config.json, the installed
|
||||
version -- and never imports a plugin module, instantiates a plugin class or
|
||||
calls a plugin lifecycle hook. This class is that read side.
|
||||
manifest, config schema, the plugin's section of config.json -- and never
|
||||
imports a plugin module, instantiates a plugin class or calls a plugin
|
||||
lifecycle hook. This class is the manifest side of that; schemas come from
|
||||
SchemaManager and config from ConfigManager.
|
||||
|
||||
It keeps the method names of the read-only part of :class:`PluginManager`
|
||||
(``discover_plugins``, ``plugin_manifests``, ``get_plugin_info``,
|
||||
@@ -25,11 +26,10 @@ machine) the web cannot know, and reports as unknown.
|
||||
See docs/ARCHITECTURE.md ("Web and display processes").
|
||||
"""
|
||||
|
||||
import json
|
||||
import threading
|
||||
import time
|
||||
from pathlib import Path
|
||||
from typing import Any, Callable, Dict, List, Optional, Union, cast
|
||||
from typing import Any, Callable, Dict, List, Optional, Union
|
||||
|
||||
from src.common.permission_utils import (
|
||||
ensure_directory_permissions, get_plugin_dir_mode,
|
||||
@@ -47,19 +47,16 @@ _RUNTIME_VIEW_TTL_SECONDS = 1.0
|
||||
|
||||
|
||||
class PluginCatalog:
|
||||
"""Manifests, schemas, config and versions of the installed plugins.
|
||||
"""Manifests and directories of the installed plugins.
|
||||
|
||||
Discovery is explicit and cheap to repeat: :meth:`discover_plugins`
|
||||
rescans the plugins directory and replaces the manifest map, so an
|
||||
uninstalled plugin disappears and a new one appears.
|
||||
"""
|
||||
|
||||
def __init__(self, plugins_dir: PathLike, config_manager: Optional[Any] = None,
|
||||
schema_manager: Optional[Any] = None,
|
||||
def __init__(self, plugins_dir: PathLike,
|
||||
runtime_source: Optional[Callable[[], Any]] = None) -> None:
|
||||
self.plugins_dir: Path = Path(plugins_dir)
|
||||
self.config_manager = config_manager
|
||||
self.schema_manager = schema_manager
|
||||
# Returns the display's PluginRuntimeView
|
||||
# (src/plugin_system/plugin_runtime.py). Its live view carries the
|
||||
# modes the display registered, which the mode lookups below prefer
|
||||
@@ -145,30 +142,6 @@ class PluginCatalog:
|
||||
ids = list(self.plugin_manifests)
|
||||
return [info for info in (self.get_plugin_info(pid) for pid in ids) if info]
|
||||
|
||||
def read_manifest(self, plugin_id: str) -> Optional[Dict[str, Any]]:
|
||||
"""The manifest as it is on disk now, not as discovery last saw it.
|
||||
|
||||
For reads that must reflect a change made since the last scan -- the
|
||||
version just after an update, say. None when the plugin has no
|
||||
directory or its manifest is missing, unreadable or not an object.
|
||||
"""
|
||||
plugin_dir = self.get_plugin_directory(plugin_id)
|
||||
if plugin_dir is None:
|
||||
return None
|
||||
try:
|
||||
with open(Path(plugin_dir) / 'manifest.json', 'r', encoding='utf-8') as f:
|
||||
manifest = json.load(f)
|
||||
except (OSError, ValueError) as exc:
|
||||
self.logger.debug("Could not read manifest for %s: %s", plugin_id, exc)
|
||||
return None
|
||||
return manifest if isinstance(manifest, dict) else None
|
||||
|
||||
def get_installed_version(self, plugin_id: str) -> str:
|
||||
"""The installed version from the on-disk manifest, or ''."""
|
||||
manifest = self.read_manifest(plugin_id) or {}
|
||||
version = manifest.get('version', '')
|
||||
return version if isinstance(version, str) else str(version)
|
||||
|
||||
def get_plugin_directory(self, plugin_id: str) -> Optional[str]:
|
||||
"""Where ``plugin_id`` is installed, or None.
|
||||
|
||||
@@ -252,31 +225,6 @@ class PluginCatalog:
|
||||
return plugin_id
|
||||
return None
|
||||
|
||||
# -- schema and config ------------------------------------------------
|
||||
|
||||
def get_schema(self, plugin_id: str, use_cache: bool = True) -> Optional[Dict[str, Any]]:
|
||||
"""The plugin's config schema through SchemaManager, or None."""
|
||||
if self.schema_manager is None:
|
||||
return None
|
||||
schema = self.schema_manager.load_schema(plugin_id, use_cache=use_cache)
|
||||
return cast(Optional[Dict[str, Any]], schema)
|
||||
|
||||
def get_config(self, plugin_id: str) -> Dict[str, Any]:
|
||||
"""The plugin's section of config.json (secrets merged), or {}."""
|
||||
if self.config_manager is None:
|
||||
return {}
|
||||
section = (self.config_manager.load_config() or {}).get(plugin_id)
|
||||
return section if isinstance(section, dict) else {}
|
||||
|
||||
def is_enabled(self, plugin_id: str) -> bool:
|
||||
"""Whether config.json enables the plugin, by the display's rule.
|
||||
|
||||
The display loads a plugin only when its section says
|
||||
``"enabled": true``; a missing flag or section means disabled
|
||||
(``DisplayController._reconcile_enabled_plugins``).
|
||||
"""
|
||||
return bool(self.get_config(plugin_id).get('enabled', False))
|
||||
|
||||
|
||||
def display_restart_required(action: str, plugin_enabled: bool, *,
|
||||
changed: bool = True,
|
||||
|
||||
@@ -38,6 +38,7 @@ from src.common.permission_utils import (
|
||||
ensure_directory_permissions,
|
||||
get_plugin_dir_mode
|
||||
)
|
||||
from src.deprecation import deprecated
|
||||
|
||||
|
||||
class _DeferredConfigChange(NamedTuple):
|
||||
@@ -939,6 +940,7 @@ class PluginManager:
|
||||
"""
|
||||
return self.plugins.get(plugin_id)
|
||||
|
||||
@deprecated("3.10.0", "use get_plugin(plugin_id)")
|
||||
def get_all_plugins(self) -> Dict[str, Any]:
|
||||
"""
|
||||
Get all loaded plugins.
|
||||
@@ -948,6 +950,7 @@ class PluginManager:
|
||||
"""
|
||||
return self.plugins.copy()
|
||||
|
||||
@deprecated("3.10.0", "read the manifest with src.plugin_system.plugin_catalog.PluginCatalog")
|
||||
def get_plugin_info(self, plugin_id: str) -> Optional[Dict[str, Any]]:
|
||||
"""
|
||||
Get information about a plugin (manifest + runtime info).
|
||||
@@ -985,6 +988,7 @@ class PluginManager:
|
||||
|
||||
return info
|
||||
|
||||
@deprecated("3.10.0", "read manifests with src.plugin_system.plugin_catalog.PluginCatalog")
|
||||
def get_all_plugin_info(self) -> List[Dict[str, Any]]:
|
||||
"""
|
||||
Get information about all plugins.
|
||||
@@ -1025,6 +1029,7 @@ class PluginManager:
|
||||
by_manifest=False)
|
||||
return str(plugin_dir) if plugin_dir is not None else None
|
||||
|
||||
@deprecated("3.10.0", "read manifests with src.plugin_system.plugin_catalog.PluginCatalog")
|
||||
def get_plugin_display_modes(self, plugin_id: str) -> List[str]:
|
||||
"""
|
||||
Get display modes provided by a plugin.
|
||||
@@ -1045,6 +1050,7 @@ class PluginManager:
|
||||
return display_modes
|
||||
return []
|
||||
|
||||
@deprecated("3.10.0", "read manifests with src.plugin_system.plugin_catalog.PluginCatalog")
|
||||
def find_plugin_for_mode(self, mode: str) -> Optional[str]:
|
||||
"""
|
||||
Find which plugin provides a given display mode.
|
||||
|
||||
@@ -15,6 +15,7 @@ from datetime import datetime
|
||||
import logging
|
||||
|
||||
from src.logging_config import get_logger
|
||||
from src.deprecation import deprecated
|
||||
|
||||
|
||||
class PluginState(Enum):
|
||||
@@ -138,6 +139,7 @@ class PluginStateManager:
|
||||
"""
|
||||
return self._states.get(plugin_id, PluginState.UNLOADED)
|
||||
|
||||
@deprecated("3.10.0", "use get_state()")
|
||||
def is_loaded(self, plugin_id: str) -> bool:
|
||||
"""Check if plugin is loaded."""
|
||||
state = self.get_state(plugin_id)
|
||||
@@ -148,11 +150,13 @@ class PluginStateManager:
|
||||
state = self.get_state(plugin_id)
|
||||
return state == PluginState.ENABLED
|
||||
|
||||
@deprecated("3.10.0", "use get_state()")
|
||||
def is_running(self, plugin_id: str) -> bool:
|
||||
"""Check if plugin is currently running."""
|
||||
state = self.get_state(plugin_id)
|
||||
return state == PluginState.RUNNING
|
||||
|
||||
@deprecated("3.10.0", "use get_state()")
|
||||
def is_error(self, plugin_id: str) -> bool:
|
||||
"""Check if plugin is in error state."""
|
||||
state = self.get_state(plugin_id)
|
||||
@@ -197,6 +201,7 @@ class PluginStateManager:
|
||||
state.value,
|
||||
)
|
||||
|
||||
@deprecated("3.10.0")
|
||||
def get_error_info(self, plugin_id: str) -> Optional[Dict[str, Any]]:
|
||||
"""
|
||||
Get error information for a plugin.
|
||||
@@ -287,10 +292,12 @@ class PluginStateManager:
|
||||
"""Record that plugin update() was called."""
|
||||
self._last_update[plugin_id] = datetime.now()
|
||||
|
||||
@deprecated("3.10.0")
|
||||
def get_last_update(self, plugin_id: str) -> Optional[datetime]:
|
||||
"""Get timestamp of last update() call."""
|
||||
return self._last_update.get(plugin_id)
|
||||
|
||||
@deprecated("3.10.0", "use get_state()")
|
||||
def get_state_info(self, plugin_id: str) -> Dict[str, Any]:
|
||||
"""
|
||||
Get comprehensive state information for a plugin.
|
||||
|
||||
@@ -38,7 +38,6 @@ class InconsistencyType(Enum):
|
||||
PLUGIN_MISSING_ON_DISK = "plugin_missing_on_disk"
|
||||
PLUGIN_ENABLED_MISMATCH = "plugin_enabled_mismatch"
|
||||
PLUGIN_VERSION_MISMATCH = "plugin_version_mismatch"
|
||||
PLUGIN_STATE_CORRUPTED = "plugin_state_corrupted"
|
||||
|
||||
|
||||
class FixAction(Enum):
|
||||
@@ -57,7 +56,6 @@ class Inconsistency:
|
||||
fix_action: FixAction
|
||||
current_state: Dict[str, Any]
|
||||
expected_state: Dict[str, Any]
|
||||
can_auto_fix: bool = False
|
||||
|
||||
|
||||
@dataclass
|
||||
@@ -270,7 +268,7 @@ class StateReconciliation:
|
||||
|
||||
# Attempt to fix auto-fixable inconsistencies
|
||||
for inconsistency in inconsistencies:
|
||||
if inconsistency.can_auto_fix and inconsistency.fix_action == FixAction.AUTO_FIX:
|
||||
if inconsistency.fix_action == FixAction.AUTO_FIX:
|
||||
if self._fix_inconsistency(inconsistency):
|
||||
fixed.append(inconsistency)
|
||||
else:
|
||||
@@ -428,7 +426,6 @@ class StateReconciliation:
|
||||
fix_action=FixAction.AUTO_FIX,
|
||||
current_state={'exists_in_config': False},
|
||||
expected_state={'exists_in_config': True, 'enabled': False},
|
||||
can_auto_fix=True
|
||||
))
|
||||
|
||||
# Check: Plugin in config but not on disk
|
||||
@@ -459,7 +456,6 @@ class StateReconciliation:
|
||||
fix_action=FixAction.AUTO_FIX if can_repair else FixAction.MANUAL_FIX_REQUIRED,
|
||||
current_state={'exists_on_disk': False},
|
||||
expected_state={'exists_on_disk': True},
|
||||
can_auto_fix=can_repair
|
||||
))
|
||||
|
||||
# Observed checks: only against a live snapshot, and only for a plugin
|
||||
@@ -486,7 +482,6 @@ class StateReconciliation:
|
||||
fix_action=FixAction.NO_ACTION,
|
||||
current_state={'loaded': loaded, 'state': runtime.get('state')},
|
||||
expected_state={'loaded': config_enabled},
|
||||
can_auto_fix=False
|
||||
))
|
||||
loaded_version = runtime.get('loaded_version')
|
||||
disk_version = disk.get('version')
|
||||
@@ -500,7 +495,6 @@ class StateReconciliation:
|
||||
fix_action=FixAction.NO_ACTION,
|
||||
current_state={'version': loaded_version},
|
||||
expected_state={'version': disk_version},
|
||||
can_auto_fix=False
|
||||
))
|
||||
|
||||
return inconsistencies
|
||||
|
||||
@@ -183,16 +183,7 @@ class _RegistryMixin:
|
||||
@staticmethod
|
||||
def _distinct_sequence(values: List[str]) -> List[str]:
|
||||
"""Return list preserving order while removing duplicates and falsey entries."""
|
||||
seen = set()
|
||||
ordered = []
|
||||
for value in values:
|
||||
if not value:
|
||||
continue
|
||||
if value in seen:
|
||||
continue
|
||||
seen.add(value)
|
||||
ordered.append(value)
|
||||
return ordered
|
||||
return list(dict.fromkeys(v for v in values if v))
|
||||
|
||||
def _validate_manifest_version_fields(self, manifest: Dict[str, Any]) -> List[str]:
|
||||
"""
|
||||
|
||||
@@ -25,6 +25,7 @@ from src.plugin_system.testing.mocks import (
|
||||
MockConfigManager,
|
||||
MockPluginManager
|
||||
)
|
||||
from src.deprecation import deprecated
|
||||
|
||||
|
||||
class PluginTestCase(unittest.TestCase):
|
||||
@@ -34,6 +35,7 @@ class PluginTestCase(unittest.TestCase):
|
||||
Provides common fixtures and helper methods.
|
||||
"""
|
||||
|
||||
@deprecated("3.10.0", "use src.plugin_system.testing.harness and the mocks directly")
|
||||
def setUp(self):
|
||||
"""Set up test fixtures."""
|
||||
# Create mock managers
|
||||
|
||||
@@ -291,49 +291,6 @@ class VegasModeConfig:
|
||||
max_cycle_duration=int(get('max_cycle_duration', d.max_cycle_duration)),
|
||||
)
|
||||
|
||||
def to_dict(self) -> Dict[str, Any]:
|
||||
"""Convert config to dictionary for serialization."""
|
||||
return {
|
||||
'enabled': self.enabled,
|
||||
'scroll_speed': self.scroll_speed,
|
||||
'separator_width': self.separator_width,
|
||||
'intra_plugin_gap': self.intra_plugin_gap,
|
||||
'render_width_pct': self.render_width_pct,
|
||||
'min_content_separation': self.min_content_separation,
|
||||
'min_cut_gap': self.min_cut_gap,
|
||||
'smooth_scroll': self.smooth_scroll,
|
||||
'sub_pixel_blend': self.sub_pixel_blend,
|
||||
'continuous_scroll': self.continuous_scroll,
|
||||
'offscreen_prefetch': self.offscreen_prefetch,
|
||||
'switch_interval_ms': self.switch_interval_ms,
|
||||
'prefetch_gate': self.prefetch_gate,
|
||||
'live_refresh': self.live_refresh,
|
||||
'live_max_hz': self.live_max_hz,
|
||||
'live_min_interval': self.live_min_interval,
|
||||
'live_lead_screens': self.live_lead_screens,
|
||||
'extend_threshold_screens': self.extend_threshold_screens,
|
||||
'auto_trim': self.auto_trim,
|
||||
'trim_threshold': self.trim_threshold,
|
||||
'content_padding': self.content_padding,
|
||||
'min_plugin_width': self.min_plugin_width,
|
||||
'lead_in_width': self.lead_in_width,
|
||||
'plugins_per_cycle': self.plugins_per_cycle,
|
||||
'max_plugin_width_ratio': self.max_plugin_width_ratio,
|
||||
'live_in_ticker': self.live_in_ticker,
|
||||
'live_weight': self.live_weight,
|
||||
'favorite_live_weight': self.favorite_live_weight,
|
||||
'overflow_mode': self.overflow_mode,
|
||||
'plugin_order': self.plugin_order,
|
||||
'excluded_plugins': list(self.excluded_plugins),
|
||||
'target_fps': self.target_fps,
|
||||
'buffer_ahead': self.buffer_ahead,
|
||||
'frame_based_scrolling': self.frame_based_scrolling,
|
||||
'scroll_delay': self.scroll_delay,
|
||||
'dynamic_duration_enabled': self.dynamic_duration_enabled,
|
||||
'min_cycle_duration': self.min_cycle_duration,
|
||||
'max_cycle_duration': self.max_cycle_duration,
|
||||
}
|
||||
|
||||
def get_frame_interval(self) -> float:
|
||||
"""Get the frame interval in seconds for target FPS."""
|
||||
return 1.0 / max(1, self.target_fps)
|
||||
|
||||
@@ -182,16 +182,6 @@ class VegasModeCoordinator:
|
||||
self._static_pause_active = False
|
||||
self._saved_scroll_position: Optional[int] = None
|
||||
|
||||
# Statistics
|
||||
self.stats = {
|
||||
'total_runtime_seconds': 0.0,
|
||||
'cycles_completed': 0,
|
||||
'interruptions': 0,
|
||||
'config_updates': 0,
|
||||
'static_pauses': 0,
|
||||
}
|
||||
self._start_time: Optional[float] = None
|
||||
|
||||
logger.info(
|
||||
"VegasModeCoordinator initialized: enabled=%s, fps=%d, buffer_ahead=%d",
|
||||
self.vegas_config.enabled,
|
||||
@@ -326,7 +316,6 @@ class VegasModeCoordinator:
|
||||
# new run would have run_frame() refuse every frame.
|
||||
self._is_paused = False
|
||||
self._live_priority_active = False
|
||||
self._start_time = time.time()
|
||||
# A fresh run starts with a clean health slate: no stale
|
||||
# "was degraded" from the previous run, and a heartbeat that is
|
||||
# due immediately so the first sample confirms the marquee is up.
|
||||
@@ -357,10 +346,6 @@ class VegasModeCoordinator:
|
||||
self._is_paused = False
|
||||
self._live_priority_active = False
|
||||
|
||||
if self._start_time:
|
||||
self.stats['total_runtime_seconds'] += time.time() - self._start_time
|
||||
self._start_time = None
|
||||
|
||||
self._restore_switch_interval()
|
||||
self._remove_render_gate()
|
||||
self._set_live(False, None)
|
||||
@@ -485,7 +470,6 @@ class VegasModeCoordinator:
|
||||
if not self._is_active:
|
||||
return
|
||||
self._is_paused = True
|
||||
self.stats['interruptions'] += 1
|
||||
|
||||
self.display_manager.set_scrolling_state(False)
|
||||
logger.info("Vegas mode paused")
|
||||
@@ -555,9 +539,8 @@ class VegasModeCoordinator:
|
||||
if self.render_pipeline.has_deferred():
|
||||
self.render_pipeline.drain_deferred()
|
||||
elif self.render_pipeline.needs_extension():
|
||||
if self.render_pipeline.extend_scroll_content():
|
||||
self.stats['cycles_completed'] += 1
|
||||
elif self.render_pipeline.is_cycle_complete():
|
||||
if (not self.render_pipeline.extend_scroll_content()
|
||||
and self.render_pipeline.is_cycle_complete()):
|
||||
# Extension failed and the strip has run out: fall back to
|
||||
# the swap rather than sitting on a dead frame.
|
||||
self.render_pipeline.start_new_cycle()
|
||||
@@ -567,7 +550,6 @@ class VegasModeCoordinator:
|
||||
if not self.render_pipeline.start_new_cycle():
|
||||
logger.warning("Failed to start new Vegas cycle")
|
||||
return False
|
||||
self.stats['cycles_completed'] += 1
|
||||
|
||||
# Check for hot-swap opportunities
|
||||
if self.render_pipeline.should_recompose():
|
||||
@@ -847,7 +829,6 @@ class VegasModeCoordinator:
|
||||
self._pending_config_update = True
|
||||
self._pending_config = new_config
|
||||
self._config_version += 1
|
||||
self.stats['config_updates'] += 1
|
||||
|
||||
logger.debug("Config update queued (version %d)", self._config_version)
|
||||
|
||||
@@ -922,23 +903,6 @@ class VegasModeCoordinator:
|
||||
self.stream_manager.mark_plugin_updated(plugin_id)
|
||||
self.plugin_adapter.invalidate_cache(plugin_id)
|
||||
|
||||
def get_status(self) -> Dict[str, Any]:
|
||||
"""Get comprehensive Vegas mode status."""
|
||||
status = {
|
||||
'enabled': self.vegas_config.enabled,
|
||||
'active': self._is_active,
|
||||
'paused': self._is_paused,
|
||||
'live_priority_active': self._live_priority_active,
|
||||
'config': self.vegas_config.to_dict(),
|
||||
'stats': self.stats.copy(),
|
||||
}
|
||||
|
||||
if self._is_active:
|
||||
status['render_info'] = self.render_pipeline.get_current_scroll_info()
|
||||
status['stream_status'] = self.stream_manager.get_buffer_status()
|
||||
|
||||
return status
|
||||
|
||||
# -------------------------------------------------------------------------
|
||||
# Static pause handling (for STATIC display mode)
|
||||
# -------------------------------------------------------------------------
|
||||
@@ -992,7 +956,6 @@ class VegasModeCoordinator:
|
||||
# Save current scroll position for smooth resume
|
||||
self._saved_scroll_position = self.render_pipeline.get_scroll_position()
|
||||
self._static_pause_active = True
|
||||
self.stats['static_pauses'] += 1
|
||||
|
||||
logger.info("Static pause started for plugin: %s", plugin_id)
|
||||
|
||||
|
||||
@@ -218,7 +218,6 @@ class RenderPipeline:
|
||||
|
||||
# Render state
|
||||
self._cycle_complete = False
|
||||
self._segments_in_scroll: List[str] = [] # Plugin IDs in current scroll
|
||||
self._record_by_seq: Dict[int, ElementRecord] = {}
|
||||
# Live updates. _applied: per record, the (epoch, digest) of the
|
||||
# pixels the strip holds. _live_slots / _live_ready: the worker's
|
||||
@@ -235,15 +234,9 @@ class RenderPipeline:
|
||||
self._frame_interval = config.get_frame_interval()
|
||||
self._cycle_start_time = 0.0
|
||||
|
||||
# Statistics
|
||||
self.stats = {
|
||||
'frames_rendered': 0,
|
||||
'scroll_cycles': 0,
|
||||
'composition_count': 0,
|
||||
'hot_swaps': 0,
|
||||
'avg_frame_time_ms': 0.0,
|
||||
}
|
||||
self._frame_times: Deque[float] = deque(maxlen=100) # Efficient fixed-size buffer
|
||||
# Read by _measure_refresh (warm-up) and the live integration test.
|
||||
self.frames_rendered = 0
|
||||
self.extensions = 0
|
||||
|
||||
logger.info(
|
||||
"RenderPipeline initialized: %dx%d @ %d FPS",
|
||||
@@ -342,7 +335,7 @@ class RenderPipeline:
|
||||
return
|
||||
if getattr(self.display_manager, 'matrix', None) is None:
|
||||
return # No hardware: nothing blocks, so there is nothing to time.
|
||||
if self.stats['frames_rendered'] < self.REFRESH_WARMUP_FRAMES:
|
||||
if self.frames_rendered < self.REFRESH_WARMUP_FRAMES:
|
||||
return
|
||||
self._swap_times.append(time.monotonic())
|
||||
if len(self._swap_times) <= self.REFRESH_SAMPLES:
|
||||
@@ -452,10 +445,6 @@ class RenderPipeline:
|
||||
layouts)
|
||||
self._note_op('compose', self._strip_nbytes())
|
||||
|
||||
# Track which plugins are in this scroll (get safely via buffer status)
|
||||
self._segments_in_scroll = self.stream_manager.get_active_plugin_ids()
|
||||
|
||||
self.stats['composition_count'] += 1
|
||||
self._cycle_start_time = time.time()
|
||||
self._cycle_complete = False
|
||||
|
||||
@@ -865,9 +854,7 @@ class RenderPipeline:
|
||||
# trim's copy, if it made one.
|
||||
self._note_op('extend', moved + (self._copied_bytes() if cut else 0))
|
||||
|
||||
self._segments_in_scroll = [pid for pid, _ in grouped]
|
||||
self.stats['composition_count'] += 1
|
||||
self.stats['extensions'] = self.stats.get('extensions', 0) + 1
|
||||
self.extensions += 1
|
||||
|
||||
logger.info(
|
||||
"Extended scroll strip with %d plugin block(s), %d rows: "
|
||||
@@ -1150,8 +1137,6 @@ class RenderPipeline:
|
||||
Returns:
|
||||
True if frame was rendered, False if no content
|
||||
"""
|
||||
frame_start = time.time()
|
||||
|
||||
try:
|
||||
if not self.scroll_helper.has_strip():
|
||||
return False
|
||||
@@ -1200,7 +1185,6 @@ class RenderPipeline:
|
||||
if at_wrap_point or self.scroll_helper.is_scroll_complete():
|
||||
if not self._cycle_complete:
|
||||
self._cycle_complete = True
|
||||
self.stats['scroll_cycles'] += 1
|
||||
logger.info(
|
||||
"Scroll cycle complete after %.1fs",
|
||||
time.time() - self._cycle_start_time
|
||||
@@ -1239,11 +1223,8 @@ class RenderPipeline:
|
||||
# Update scrolling state
|
||||
self.display_manager.set_scrolling_state(True, self._frame_hold)
|
||||
|
||||
# Track statistics
|
||||
self.stats['frames_rendered'] += 1
|
||||
self.frames_rendered += 1
|
||||
self._measure_refresh()
|
||||
frame_time = time.time() - frame_start
|
||||
self._track_frame_time(frame_time)
|
||||
|
||||
return True
|
||||
|
||||
@@ -1252,15 +1233,6 @@ class RenderPipeline:
|
||||
logger.exception("Error rendering frame")
|
||||
return False
|
||||
|
||||
def _track_frame_time(self, frame_time: float) -> None:
|
||||
"""Track frame timing for statistics."""
|
||||
self._frame_times.append(frame_time) # deque with maxlen auto-removes old entries
|
||||
|
||||
if self._frame_times:
|
||||
self.stats['avg_frame_time_ms'] = (
|
||||
sum(self._frame_times) / len(self._frame_times) * 1000
|
||||
)
|
||||
|
||||
def is_cycle_complete(self) -> bool:
|
||||
"""Check if current scroll cycle is complete."""
|
||||
return self._cycle_complete
|
||||
@@ -1347,7 +1319,6 @@ class RenderPipeline:
|
||||
else:
|
||||
self.scroll_helper.scroll_position = 0.0
|
||||
|
||||
self.stats['hot_swaps'] += 1
|
||||
logger.debug(
|
||||
"Hot-swap completed: scroll repositioned %.0f→%.0f (%.1f%% of new %dpx image)",
|
||||
old_pos, self.scroll_helper.scroll_position,
|
||||
@@ -1396,8 +1367,6 @@ class RenderPipeline:
|
||||
# transition rather than near-end content wrapping around.
|
||||
self.scroll_helper.scroll_position = float(self.config.lead_in_width)
|
||||
|
||||
# Signal follower that a new cycle started (triggers its own rebuild)
|
||||
self.sync_manager.send_new_cycle()
|
||||
# Push the actual scroll image over TCP so follower has identical pixels.
|
||||
# Done in a background thread to not block the render loop (~15ms transfer).
|
||||
image = self.scroll_helper.cached_image
|
||||
@@ -1410,16 +1379,6 @@ class RenderPipeline:
|
||||
|
||||
return result
|
||||
|
||||
def get_current_scroll_info(self) -> Dict[str, Any]:
|
||||
"""Get current scroll state information."""
|
||||
scroll_info = self.scroll_helper.get_scroll_info()
|
||||
return {
|
||||
**scroll_info,
|
||||
'cycle_complete': self._cycle_complete,
|
||||
'plugins_in_scroll': self._segments_in_scroll,
|
||||
'stats': self.stats.copy(),
|
||||
}
|
||||
|
||||
def get_scroll_position(self) -> int:
|
||||
"""
|
||||
Get current scroll position.
|
||||
@@ -1465,8 +1424,6 @@ class RenderPipeline:
|
||||
self.scroll_helper.clear_cache()
|
||||
|
||||
self._cycle_complete = False
|
||||
self._segments_in_scroll = []
|
||||
self._frame_times = deque(maxlen=100)
|
||||
|
||||
# Content lined up for the old run belongs to it. Left in place, the
|
||||
# first extension after Vegas is switched back on appended that stale
|
||||
|
||||
@@ -16,7 +16,7 @@ BasePlugin.get_vegas_participation):
|
||||
import logging
|
||||
import threading
|
||||
import time
|
||||
from typing import Optional, List, Dict, Any, Deque, Tuple, TYPE_CHECKING
|
||||
from typing import Optional, List, Dict, Deque, Tuple, TYPE_CHECKING
|
||||
from collections import deque
|
||||
from dataclasses import dataclass, field
|
||||
from PIL import Image
|
||||
@@ -94,13 +94,6 @@ class StreamManager:
|
||||
self._last_refresh: float = 0.0
|
||||
self._refresh_interval: float = 30.0 # Refresh plugin list every 30s
|
||||
|
||||
# Statistics
|
||||
self.stats = {
|
||||
'segments_fetched': 0,
|
||||
'segments_served': 0,
|
||||
'fetch_errors': 0,
|
||||
}
|
||||
|
||||
logger.info("StreamManager initialized with buffer_ahead=%d", config.buffer_ahead)
|
||||
|
||||
def initialize(self) -> bool:
|
||||
@@ -143,47 +136,12 @@ class StreamManager:
|
||||
return None
|
||||
|
||||
segment = self._active_buffer.popleft()
|
||||
self.stats['segments_served'] += 1
|
||||
|
||||
# Trigger prefetch to maintain buffer
|
||||
self._ensure_buffer_filled()
|
||||
|
||||
return segment
|
||||
|
||||
def peek_next_segment(self) -> Optional[ContentSegment]:
|
||||
"""
|
||||
Peek at the next segment without removing it.
|
||||
|
||||
Returns:
|
||||
ContentSegment or None if buffer is empty
|
||||
"""
|
||||
with self._buffer_lock:
|
||||
if self._active_buffer:
|
||||
return self._active_buffer[0]
|
||||
return None
|
||||
|
||||
def get_buffer_status(self) -> Dict[str, Any]:
|
||||
"""Get current buffer status for monitoring."""
|
||||
with self._buffer_lock:
|
||||
return {
|
||||
'active_count': len(self._active_buffer),
|
||||
'total_plugins': len(self._ordered_plugins),
|
||||
'prefetch_index': self._prefetch_index,
|
||||
'stats': self.stats.copy(),
|
||||
}
|
||||
|
||||
def get_active_plugin_ids(self) -> List[str]:
|
||||
"""
|
||||
Get list of plugin IDs currently in the active buffer.
|
||||
|
||||
Thread-safe accessor for render pipeline.
|
||||
|
||||
Returns:
|
||||
List of plugin IDs in buffer order
|
||||
"""
|
||||
with self._buffer_lock:
|
||||
return [seg.plugin_id for seg in self._active_buffer]
|
||||
|
||||
def mark_plugin_updated(self, plugin_id: str) -> None:
|
||||
"""
|
||||
Mark a plugin as having updated data.
|
||||
@@ -587,7 +545,6 @@ class StreamManager:
|
||||
images=[], # No images needed for static pause
|
||||
display_mode=VegasDisplayMode.STATIC
|
||||
)
|
||||
self.stats['segments_fetched'] += 1
|
||||
logger.debug(
|
||||
"[%s] Created STATIC placeholder (pause trigger)",
|
||||
plugin_id
|
||||
@@ -610,7 +567,6 @@ class StreamManager:
|
||||
display_mode=VegasDisplayMode.SCROLL
|
||||
)
|
||||
|
||||
self.stats['segments_fetched'] += 1
|
||||
logger.debug(
|
||||
"[%s] Segment: %d image(s), %dpx",
|
||||
plugin_id, len(images), total_width
|
||||
@@ -619,7 +575,6 @@ class StreamManager:
|
||||
|
||||
except Exception:
|
||||
logger.exception("[%s] ERROR fetching content", plugin_id)
|
||||
self.stats['fetch_errors'] += 1
|
||||
return None
|
||||
|
||||
def _ensure_buffer_filled(self) -> None:
|
||||
@@ -770,10 +725,8 @@ class StreamManager:
|
||||
plugin, plugin_id, offscreen_only=offscreen_only)
|
||||
except Exception:
|
||||
logger.exception("[%s] ERROR fetching content", plugin_id)
|
||||
self.stats['fetch_errors'] += 1
|
||||
return None
|
||||
if images:
|
||||
self.stats['segments_fetched'] += 1
|
||||
return (plugin_id, images)
|
||||
# Only the old contract hands anything back to the render thread.
|
||||
defer_empty = offscreen_only and not getattr(
|
||||
|
||||
@@ -8,7 +8,6 @@ import time
|
||||
from typing import Any, Optional, Dict, Tuple
|
||||
from flask import jsonify, request
|
||||
|
||||
from src.web_interface.error_handler import create_error_response, create_success_response
|
||||
from src.web_interface.errors import ErrorCode, WebInterfaceError
|
||||
|
||||
|
||||
@@ -31,7 +30,15 @@ def success_response(
|
||||
Returns:
|
||||
Flask jsonify response
|
||||
"""
|
||||
response_data = create_success_response(data, message, metadata)
|
||||
response_data: Dict[str, Any] = {'status': 'success'}
|
||||
# `is not None` rather than truthiness: "" and {} are values a caller
|
||||
# chose to send, and dropping them would make the shape depend on the data.
|
||||
if data is not None:
|
||||
response_data['data'] = data
|
||||
if message is not None:
|
||||
response_data['message'] = message
|
||||
if metadata is not None:
|
||||
response_data['metadata'] = metadata
|
||||
for key, value in (extra or {}).items():
|
||||
response_data.setdefault(key, value)
|
||||
|
||||
@@ -70,14 +77,14 @@ def error_response(
|
||||
Returns:
|
||||
Flask jsonify response with status code
|
||||
"""
|
||||
return create_error_response(
|
||||
error = WebInterfaceError(
|
||||
error_code=error_code,
|
||||
message=message,
|
||||
details=details,
|
||||
context=context,
|
||||
suggested_fixes=suggested_fixes,
|
||||
status_code=status_code
|
||||
context=context or {},
|
||||
suggested_fixes=suggested_fixes
|
||||
)
|
||||
return jsonify(error.to_dict()), status_code
|
||||
|
||||
|
||||
def exception_error_response(
|
||||
|
||||
@@ -12,8 +12,8 @@ from typing import Any, Dict
|
||||
def _schema_type_is(prop: Any, wanted: str) -> bool:
|
||||
"""Whether a schema property is of ``wanted`` type, unions included.
|
||||
|
||||
Mirrors ``_schema_type_is`` in ``web_interface/blueprints/api_v3`` (kept
|
||||
here so src/ doesn't import the Flask blueprint). A union such as
|
||||
The one copy: ``web_interface/blueprints/api_v3`` imports it from here
|
||||
(src/ must not import the Flask blueprint). A union such as
|
||||
``["array", "null"]`` -- the per-element style overrides, where null means
|
||||
"inherit" -- is still an array for recombining position-keyed inputs.
|
||||
"""
|
||||
|
||||
@@ -1,20 +1,13 @@
|
||||
"""
|
||||
Centralized error handling for web interface.
|
||||
Error text and payloads for web interface responses.
|
||||
|
||||
Provides helpers for consistent error responses across API endpoints.
|
||||
Safe exception descriptions and the bodies for exceptions no route handled.
|
||||
The standard success/error responses are in api_helpers.
|
||||
"""
|
||||
|
||||
from typing import Any, Optional
|
||||
from flask import jsonify
|
||||
|
||||
from src.web_interface.errors import WebInterfaceError, ErrorCode
|
||||
from src.logging_config import get_logger
|
||||
from src.redaction import redact_credentials
|
||||
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
|
||||
# Long enough for an errno string with a path, short enough not to dump a
|
||||
# parser's worth of context into a JSON field.
|
||||
_MAX_DETAIL_LENGTH = 400
|
||||
@@ -101,72 +94,3 @@ def http_exception_payload(error) -> dict:
|
||||
'error_code': (error.name or 'HTTP_ERROR').upper().replace(' ', '_'),
|
||||
'message': error.description,
|
||||
}
|
||||
|
||||
|
||||
def create_error_response(
|
||||
error_code: ErrorCode,
|
||||
message: str,
|
||||
details: Optional[str] = None,
|
||||
context: Optional[dict] = None,
|
||||
suggested_fixes: Optional[list] = None,
|
||||
status_code: int = 500
|
||||
) -> tuple:
|
||||
"""
|
||||
Create a standardized error response.
|
||||
|
||||
Args:
|
||||
error_code: Error code
|
||||
message: Error message
|
||||
details: Optional detailed error information
|
||||
context: Optional context dictionary
|
||||
suggested_fixes: Optional list of suggested fixes
|
||||
status_code: HTTP status code
|
||||
|
||||
Returns:
|
||||
Tuple of (jsonify response, status_code)
|
||||
"""
|
||||
error = WebInterfaceError(
|
||||
error_code=error_code,
|
||||
message=message,
|
||||
details=details,
|
||||
context=context or {},
|
||||
suggested_fixes=suggested_fixes
|
||||
)
|
||||
|
||||
return jsonify(error.to_dict()), status_code
|
||||
|
||||
|
||||
def create_success_response(
|
||||
data: Any = None,
|
||||
message: Optional[str] = None,
|
||||
metadata: Optional[dict] = None
|
||||
) -> dict:
|
||||
"""
|
||||
Create a standardized success response.
|
||||
|
||||
Args:
|
||||
data: Response data
|
||||
message: Optional success message
|
||||
metadata: Optional metadata (timing, version, etc.)
|
||||
|
||||
Returns:
|
||||
Dictionary for jsonify
|
||||
"""
|
||||
response: dict[str, Any] = {
|
||||
"status": "success"
|
||||
}
|
||||
|
||||
# All three use `is not None` rather than truthiness: "" and {} are
|
||||
# values a caller chose to send, and dropping them silently would make
|
||||
# the response shape depend on the data.
|
||||
if data is not None:
|
||||
response["data"] = data
|
||||
|
||||
if message is not None:
|
||||
response["message"] = message
|
||||
|
||||
if metadata is not None:
|
||||
response["metadata"] = metadata
|
||||
|
||||
return response
|
||||
|
||||
|
||||
@@ -61,27 +61,14 @@ class WebInterfaceError:
|
||||
context: Optional[Dict[str, Any]] = None
|
||||
suggested_fixes: Optional[List[str]] = None
|
||||
original_error: Optional[Exception] = None
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
error_code: ErrorCode,
|
||||
message: str,
|
||||
details: Optional[str] = None,
|
||||
context: Optional[Dict[str, Any]] = None,
|
||||
suggested_fixes: Optional[List[str]] = None,
|
||||
original_error: Optional[Exception] = None
|
||||
):
|
||||
self.error_code = error_code
|
||||
self.message = message
|
||||
self.details = details
|
||||
self.context = context or {}
|
||||
|
||||
def __post_init__(self) -> None:
|
||||
self.context = self.context or {}
|
||||
# `is None`, not truthiness: an explicit [] means "this caller has
|
||||
# no suggestions to offer", which the default list would override.
|
||||
self.suggested_fixes = (
|
||||
suggested_fixes if suggested_fixes is not None
|
||||
else self._get_default_suggestions(error_code))
|
||||
self.original_error = original_error
|
||||
|
||||
if self.suggested_fixes is None:
|
||||
self.suggested_fixes = self._get_default_suggestions(self.error_code)
|
||||
|
||||
def _get_default_suggestions(self, error_code: ErrorCode) -> List[str]:
|
||||
"""Get default suggested fixes for error code."""
|
||||
suggestions_map = {
|
||||
|
||||
Reference in New Issue
Block a user