from flask import Flask, request, redirect, url_for, jsonify, Response, send_from_directory import json import logging import os import queue import re import gzip import shutil import sys import subprocess import threading import time from pathlib import Path from datetime import datetime, timedelta # Add parent directory to path for imports sys.path.insert(0, str(Path(__file__).parent.parent)) # Configure logging before anything below logs: the same setup as the display # service (run.py), so this process's journal lines carry their real syslog # priority too (`journalctl -p err -u ledmatrix-web`). LEDMATRIX_DEBUG=true # turns on DEBUG, which includes the routine per-request lines. from src.logging_config import setup_logging setup_logging(format_type=( 'json' if os.environ.get('LEDMATRIX_JSON_LOGGING', 'false').lower() == 'true' else 'readable')) logging.getLogger('werkzeug').setLevel(logging.WARNING) # request_logging covers requests logging.getLogger('urllib3').setLevel(logging.WARNING) from src.config_manager import ConfigManager from src.web_interface.error_handler import describe_exception from src.common.path_safety import ( resolve_under, safe_path_component, safe_relative_parts, ) from werkzeug.exceptions import HTTPException from src.exceptions import ConfigError from src.plugin_system.plugin_catalog import PluginCatalog from src.plugin_system.store_manager import PluginStoreManager from src.plugin_system.saved_repositories import SavedRepositoriesManager from src.plugin_system.schema_manager import SchemaManager from src.plugin_system.operation_queue import PluginOperationQueue from src.plugin_system.operation_history import OperationHistory _JOURNALCTL = shutil.which('journalctl') _SYSTEMCTL = shutil.which('systemctl') _NMCLI = shutil.which('nmcli') _VCGENCMD = shutil.which('vcgencmd') from web_interface import display_preview from web_interface.system_metrics import collect_system_metrics # Create Flask app app = Flask(__name__) app.secret_key = os.urandom(24) config_manager = ConfigManager() # Cross-site request forgery: the UI has no login, and being "only on the LAN" # does not keep other websites out. Any page a LAN user opens can make their # browser POST to this server -- a plain HTML form is not blocked by CORS -- so # a hostile site could reboot the Pi, pull code or rewrite the config through # the user's browser. web_interface/origin_guard.py (registered below) refuses # POST/PUT/PATCH/DELETE whose Origin (or, failing that, Referer) is not this # server's own host; requests with neither header (curl, Home Assistant, the # MQTT bridge) are not from a browser and pass. There are no CSRF tokens: # neither the HTMX forms nor the fetch() calls carry one. Anyone who can reach # the port directly can still use the API unless the optional login is on: # web_interface/auth.py (registered below the captive-portal redirect) adds a # password and API tokens. It is off until a password is set in General > # Security, and even then leaves requests from the Pi itself and the Wi-Fi # setup flow in access-point mode open. # Initialize rate limiting (prevent accidental abuse, not security) try: from flask_limiter import Limiter from flask_limiter.util import get_remote_address limiter = Limiter( app=app, key_func=get_remote_address, default_limits=["1000 per minute"], # Generous limit for local use storage_uri="memory://" # In-memory storage for simplicity ) except ImportError: # flask-limiter not installed, rate limiting disabled limiter = None # Enable gzip/brotli response compression (Flask-Compress skips streaming # responses, so the SSE endpoints are unaffected). Optional, like limiter: # missing package just means uncompressed responses. _HAVE_FLASK_COMPRESS = False try: from flask_compress import Compress Compress(app) _HAVE_FLASK_COMPRESS = True except ImportError: logging.getLogger(__name__).warning( "flask-compress not installed - responses will be served uncompressed. " "Install it with the Tools tab's 'Install Base Requirements' button or " "'pip install flask-compress'." ) # Initialize plugin managers - read plugins directory from config config = config_manager.load_config() plugin_system_config = config.get('plugin_system', {}) plugins_dir_name = plugin_system_config.get('plugins_directory', 'plugin-repos') # Project root (LEDMatrix directory). Needed below for data/ and assets/ paths # whether or not the plugins directory is absolute. project_root = Path(__file__).parent.parent # Resolve plugin directory - handle both absolute and relative paths if os.path.isabs(plugins_dir_name): plugins_dir = Path(plugins_dir_name) else: # If relative, resolve relative to the project root plugins_dir = project_root / plugins_dir_name plugin_store_manager = PluginStoreManager(plugins_dir=str(plugins_dir)) # A core `git pull` update (or any checkout) restores built-in plugins # committed under plugin-repos/, even ones the user uninstalled. Re-remove any # the user previously uninstalled at startup so a manual update on the Pi # doesn't resurrect them. try: _purged = plugin_store_manager.purge_uninstalled_plugins() if _purged: logging.getLogger(__name__).info( "Re-removed %d uninstalled plugin(s) restored since last run: %s", len(_purged), ", ".join(_purged), ) except (OSError, RuntimeError) as _purge_err: logging.getLogger(__name__).warning( "Startup plugin purge failed: %s", _purge_err ) saved_repositories_manager = SavedRepositoriesManager() # Initialize schema manager schema_manager = SchemaManager( plugins_dir=plugins_dir, project_root=project_root, logger=None, config_manager=config_manager ) # The web process reads plugins as files and never runs them: no plugin module # is imported, no plugin class instantiated, no lifecycle hook called here. # Only the display process (src/display_controller.py) does that. Config # saves reach the running plugins through the display's config watcher; what # the display knows at run time (health, metrics, errors, current mode) it # publishes to the shared cache. See docs/ARCHITECTURE.md. plugin_catalog = PluginCatalog( plugins_dir=plugins_dir, config_manager=config_manager, schema_manager=schema_manager, ) # Initialize operation queue for plugin operations operation_queue = PluginOperationQueue(max_history=500) # No plugin state file: data/plugin_state.json is retired. Desired state is # config.json plus the plugins on disk, observed state is the runtime # snapshot the display publishes (src/plugin_system/plugin_runtime.py). An # existing file is left where it is, unread; see docs/ARCHITECTURE.md. # Initialize operation history # Use lazy_load=True to defer file loading until first use (improves startup time) operation_history = OperationHistory( history_file=str(project_root / "data" / "operation_history.json"), max_records=1000, lazy_load=True ) # Plugin discovery is deferred until first API request that needs it # This improves startup time - endpoints call plugin_catalog.discover_plugins() when needed # Register blueprints from web_interface.blueprints.pages_v3 import pages_v3 from web_interface.blueprints.api_v3 import api_v3 # Initialize managers in blueprints pages_v3.config_manager = config_manager pages_v3.plugin_catalog = plugin_catalog pages_v3.plugin_store_manager = plugin_store_manager pages_v3.saved_repositories_manager = saved_repositories_manager pages_v3.schema_manager = schema_manager api_v3.config_manager = config_manager api_v3.plugin_catalog = plugin_catalog api_v3.plugin_store_manager = plugin_store_manager api_v3.saved_repositories_manager = saved_repositories_manager api_v3.schema_manager = schema_manager api_v3.operation_queue = operation_queue api_v3.operation_history = operation_history # Initialize cache manager for API endpoints from src.cache_manager import CacheManager api_v3.cache_manager = CacheManager() # Plugin health and metrics as the display publishes them. The display service # records health and execution-time metrics to the shared on-disk cache; a # tracker/monitor backed by that same cache lets the health API routes # (/api/v3/plugins/health, /plugins/metrics) read what it wrote. # Guarded so any init failure degrades to "not available" rather than breaking # the web server. api_v3.health_tracker = None api_v3.resource_monitor = None try: from src.plugin_system.plugin_health import PluginHealthTracker from src.plugin_system.resource_monitor import PluginResourceMonitor api_v3.health_tracker = PluginHealthTracker(api_v3.cache_manager) api_v3.resource_monitor = PluginResourceMonitor(api_v3.cache_manager) except Exception as _hm_err: # pragma: no cover - defensive startup guard logging.getLogger(__name__).warning( "Could not enable plugin health/metrics for web UI: %s", _hm_err ) # Pages are served un-prefixed (the interface lives at /); the /v3 mount is a # legacy alias kept so existing bookmarks and the hardcoded /v3/partials/... # fetches in templates/JS keep working unchanged. url_for('pages_v3.*') # resolves against the primary (un-prefixed) registration. app.register_blueprint(pages_v3, url_prefix='') app.register_blueprint(pages_v3, url_prefix='/v3', name='pages_v3_legacy') app.register_blueprint(api_v3, url_prefix='/api/v3') # Route to serve plugin asset files (registered on main app, not blueprint, for /assets/... path) @app.route('/assets/plugins//uploads/', methods=['GET']) def serve_plugin_asset(plugin_id, filename): """Serve uploaded asset files from assets/plugins/{plugin_id}/uploads/ Both URL parts are validated before any path is built. The containment check here used to compare the *asset directory* against the project root rather than against assets/plugins, so a plugin_id of ``..`` moved the served directory a level up and still passed. """ try: uploads_base = (project_root / 'assets' / 'plugins').resolve() safe_plugin_id = safe_path_component(plugin_id) if not safe_plugin_id: return jsonify({'status': 'error', 'message': 'Invalid asset path'}), 403 safe_parts = safe_relative_parts(filename) if not safe_parts: return jsonify({'status': 'error', 'message': 'Invalid file path'}), 403 assets_dir = resolve_under(uploads_base, safe_plugin_id, 'uploads') if assets_dir is None: return jsonify({'status': 'error', 'message': 'Invalid asset path'}), 403 if not assets_dir.exists() or not assets_dir.is_dir(): return jsonify({'status': 'error', 'message': 'Asset directory not found'}), 404 # Resolve the requested file path. resolve_under repeats the # containment check after resolving, which is what catches a symlink # inside the uploads directory pointing out of it. requested_file = resolve_under(assets_dir, *safe_parts) if requested_file is None: return jsonify({'status': 'error', 'message': 'Invalid file path'}), 403 # Check if file exists if not requested_file.exists() or not requested_file.is_file(): return jsonify({'status': 'error', 'message': 'File not found'}), 404 # Determine content type based on file extension lowered = requested_file.name.lower() content_type = 'application/octet-stream' if lowered.endswith(('.png', '.jpg', '.jpeg')): content_type = 'image/jpeg' if lowered.endswith(('.jpg', '.jpeg')) else 'image/png' elif lowered.endswith('.gif'): content_type = 'image/gif' elif lowered.endswith('.bmp'): content_type = 'image/bmp' elif lowered.endswith('.webp'): content_type = 'image/webp' elif lowered.endswith('.svg'): content_type = 'image/svg+xml' elif lowered.endswith('.json'): content_type = 'application/json' elif lowered.endswith('.txt'): content_type = 'text/plain' # Use send_from_directory to serve the file. The path handed over is # the validated one, rebuilt from the components that were checked. return send_from_directory( str(assets_dir), '/'.join(safe_parts), mimetype=content_type ) except Exception: app.logger.exception('Error serving plugin asset file') return jsonify({ 'status': 'error', 'message': 'Internal server error' }), 500 # Prime psutil CPU measurement once at startup so interval=None returns a real value try: import psutil as _psutil_prime _psutil_prime.cpu_percent(interval=None) except ImportError: pass # systemctl answers, memoised so they are not a subprocess fork per request # (AP mode) or per SSE tick (display service). A failed check keeps the last # known answer for the same TTL rather than retrying on every request. from web_interface.cache import TTLCache from src.wifi_manager import AP_PROFILE_NAME _service_status_cache = TTLCache() _AP_MODE_CACHE_TTL = 30 # seconds — AP mode is user-initiated; 30s is fine _LEDMATRIX_SERVICE_CACHE_TTL = 15 # seconds # The only units _unit_is_active() may ask systemctl about: its argv is built # from these literals, never from request data. _CHECKABLE_UNITS = frozenset({'hostapd', 'ledmatrix'}) def _unit_is_active(unit, ttl): """`systemctl is-active `, cached for ``ttl`` seconds. False where there is no systemctl (a dev machine); on a failed check, the last known answer. """ if unit not in _CHECKABLE_UNITS: raise ValueError(f"not a checkable unit: {unit!r}") active = _service_status_cache.get(unit) if active is not None: return active active = _service_status_cache.peek(unit, False) if _SYSTEMCTL: try: result = subprocess.run([_SYSTEMCTL, 'is-active', unit], # nosec B603 - list argv, unit is from _CHECKABLE_UNITS # nosemgrep capture_output=True, text=True, timeout=2) active = result.stdout.strip() == 'active' except (subprocess.SubprocessError, OSError) as e: logging.getLogger('web_interface').warning( "systemctl is-active %s failed: %s", unit, e) _service_status_cache.set(unit, active, ttl=ttl) return active def _nmcli_ap_is_active(ttl): """Whether NetworkManager has our access-point profile up, cached for ``ttl`` seconds. WiFiManager.enable_ap_mode falls back to an nmcli AP when hostapd is not available, and there is no hostapd unit running then. The match is WiFiManager._get_ap_status_nmcli's: our profile name, or a connection of a hotspot type. False where there is no nmcli; on a failed check, the last known answer. """ active = _service_status_cache.get('nmcli-ap') if active is not None: return active active = _service_status_cache.peek('nmcli-ap', False) if _NMCLI: try: result = subprocess.run( # nosec B603 - fixed argv # nosemgrep [_NMCLI, '-t', '-f', 'NAME,TYPE', 'connection', 'show', '--active'], capture_output=True, text=True, timeout=2) active = False for line in result.stdout.splitlines(): parts = line.split(':') if len(parts) < 2: continue if parts[0].strip() == AP_PROFILE_NAME or 'hotspot' in parts[1].strip().lower(): active = True break except (subprocess.SubprocessError, OSError) as e: logging.getLogger('web_interface').warning( "nmcli active-connection check failed: %s", e) _service_status_cache.set('nmcli-ap', active, ttl=ttl) return active def is_ap_mode_active(): """ Check if access point mode is currently active (cached, 30s TTL). Uses direct systemctl/nmcli checks instead of instantiating WiFiManager, and the same two WiFiManager._is_ap_mode_active uses: hostapd, else the nmcli AP it falls back to. Checking hostapd alone meant the captive portal never triggered on a Pi whose AP came up through nmcli. """ return (_unit_is_active('hostapd', _AP_MODE_CACHE_TTL) or _nmcli_ap_is_active(_AP_MODE_CACHE_TTL)) # Captive portal detection endpoints # When AP mode is active, return responses that TRIGGER the captive portal popup. # When not in AP mode, return normal "success" responses so connectivity checks pass. @app.route('/hotspot-detect.html') def hotspot_detect(): """iOS/macOS captive portal detection endpoint""" if is_ap_mode_active(): # Non-"Success" title triggers iOS captive portal popup return redirect(url_for('pages_v3.captive_setup'), code=302) return 'SuccessSuccess', 200 @app.route('/generate_204') def generate_204(): """Android captive portal detection endpoint""" if is_ap_mode_active(): # Android expects 204 = "internet works". Non-204 triggers portal popup. return redirect(url_for('pages_v3.captive_setup'), code=302) return '', 204 @app.route('/connecttest.txt') def connecttest_txt(): """Windows captive portal detection endpoint""" if is_ap_mode_active(): return redirect(url_for('pages_v3.captive_setup'), code=302) return 'Microsoft Connect Test', 200 @app.route('/success.txt') def success_txt(): """Firefox captive portal detection endpoint""" if is_ap_mode_active(): return redirect(url_for('pages_v3.captive_setup'), code=302) return 'success', 200 # Request timing and logging (routine reads at DEBUG; see request_logging) from web_interface import request_logging request_logging.init_app(app) # Refuse state-changing requests sent by another website's page (see the # cross-site note near the top of this file). from web_interface import origin_guard origin_guard.init_app(app) # Global error handlers @app.errorhandler(404) def not_found_error(error): """Handle 404 errors.""" return jsonify({ 'status': 'error', 'error_code': 'NOT_FOUND', 'message': 'Resource not found', 'path': request.path }), 404 @app.errorhandler(500) def internal_error(error): """Handle 500 errors.""" logger = logging.getLogger('web_interface') logger.error("Internal server error", exc_info=True) payload = { 'status': 'error', 'error_code': 'INTERNAL_ERROR', 'message': 'An internal error occurred; see logs for details', } # Flask hands the original exception over as `error.original_exception` # when propagation is off; without it there is nothing to describe. original = getattr(error, 'original_exception', None) or ( error if isinstance(error, BaseException) else None) if original is not None: payload['details'] = describe_exception(original) return jsonify(payload), 500 @app.errorhandler(Exception) def handle_exception(error): """Handle all unhandled exceptions. Returning only "see logs for details" is fine until the logs are exactly what you cannot reach. A device with failing storage answered every endpoint with that sentence -- including the log viewer, because journalctl could not be executed -- while the exception underneath said `[Errno 5] Input/output error`. Naming the error costs nothing here and is frequently the whole diagnosis, so include it alongside the log pointer. """ # Werkzeug's HTTPExceptions subclass Exception, so this catch-all sees # them too and was reporting every 405, 400, 413 and 415 as a server-side # UNKNOWN_ERROR 500. A GET on a POST-only route came back as "an error # occurred" rather than "method not allowed", which tells the caller # nothing and blames the wrong side. Hand those back as themselves. if isinstance(error, HTTPException): return jsonify({ 'status': 'error', 'error_code': (error.name or 'HTTP_ERROR').upper().replace(' ', '_'), 'message': error.description, }), error.code or 500 logger = logging.getLogger('web_interface') logger.error("Unhandled exception", exc_info=True) return jsonify({ 'status': 'error', 'error_code': 'UNKNOWN_ERROR', 'message': 'An error occurred; see logs for details', 'details': describe_exception(error), }), 500 # Captive portal redirect middleware @app.before_request def captive_portal_redirect(): """ Redirect all HTTP requests to WiFi setup page when AP mode is active. This creates a captive portal experience where users are automatically directed to the WiFi configuration page. """ # Check if AP mode is active if not is_ap_mode_active(): return None # Continue normal request processing # Get the request path path = request.path # List of paths that should NOT be redirected (allow normal operation) allowed_paths = [ '/v3', # Legacy-prefixed interface and all sub-paths '/setup', # Captive setup page itself (un-prefixed mount) '/partials/', # HTMX partials (un-prefixed mount) '/settings/', # Settings search index (un-prefixed mount) '/plugin-ui/', # Plugin-provided web UI assets (un-prefixed mount) '/api/v3/', # All API endpoints '/static/', # Static files (CSS, JS, images) '/hotspot-detect.html', # iOS/macOS detection '/generate_204', # Android detection '/connecttest.txt', # Windows detection '/success.txt', # Firefox detection '/favicon.ico', # Favicon '/login', # Optional web login (web_interface/auth.py) '/logout', ] for allowed_path in allowed_paths: if path.startswith(allowed_path): return None # Redirect to lightweight captive portal setup page (not the full UI) return redirect(url_for('pages_v3.captive_setup'), code=302) # Optional login (off until a password is set in General > Security). After # the captive-portal redirect, so in AP mode an unknown path still lands on # /setup rather than on the login page; the setup flow itself stays open. from web_interface import auth as web_auth web_auth.init_app(app, config_manager, limiter=limiter, is_ap_mode_active=is_ap_mode_active) # Append a content-version query param (file mtime) to every static URL so the # long-lived `immutable` cache (see add_security_headers below) is actually safe: # when a static file changes its URL changes, so browsers refetch it. Without # this, edited JS/CSS were served immutable under an unchanging URL and never # reached clients until a manual cache clear. @app.url_defaults def add_static_version(endpoint, values): if endpoint == 'static' and values.get('filename'): try: file_path = os.path.join(app.static_folder, values['filename']) values['v'] = int(os.path.getmtime(file_path)) except OSError: # File missing (e.g. plugin asset not yet installed) — skip versioning. pass # Gzip fallback for when flask-compress isn't installed. Without it the UI # ships ~1.2 MB of uncompressed JS to a phone over WiFi. Compressed bytes are # cached per URL+version, so the Pi compresses each asset once, not per request. # URL-versioned assets: safe to cache as immutable (and to gzip once). # /assets/ is the widget bundle from pages_v3 (also reachable under /v3). _VERSIONED_ASSET_PREFIXES = ('/static/', '/assets/', '/v3/assets/') _GZIP_MIN_BYTES = 1024 _GZIP_TYPES = ( 'text/html', 'text/css', 'text/plain', 'text/javascript', 'application/javascript', 'application/json', 'image/svg+xml', ) _GZIP_CACHE_MAX_BYTES = 4 * 1024 * 1024 _gzip_cache = {} _gzip_cache_bytes = 0 _gzip_cache_lock = threading.Lock() @app.after_request def compress_text_responses(response): """Gzip text responses when flask-compress is absent (SSE untouched).""" if _HAVE_FLASK_COMPRESS or response.status_code != 200: return response if 'Content-Encoding' in response.headers: return response # send_file responses report is_streamed (they wrap a file) but have a # known size; genuinely streamed bodies (SSE, generators) are left alone. # SSE is also excluded by its text/event-stream mimetype below. if response.is_streamed and not response.direct_passthrough: return response mimetype = (response.mimetype or '').lower() if mimetype not in _GZIP_TYPES: return response if 'gzip' not in (request.headers.get('Accept-Encoding') or '').lower(): return response cache_key = None if request.path.startswith(_VERSIONED_ASSET_PREFIXES): cache_key = (request.full_path, response.headers.get('Last-Modified')) with _gzip_cache_lock: cached = _gzip_cache.get(cache_key) if cached is not None: return _apply_gzip(response, cached) # send_file responses stream from disk; materialize before compressing. if response.direct_passthrough: response.direct_passthrough = False data = response.get_data() if len(data) < _GZIP_MIN_BYTES: return response compressed = gzip.compress(data, compresslevel=6) if len(compressed) >= len(data): return response if cache_key is not None: global _gzip_cache_bytes with _gzip_cache_lock: if _gzip_cache_bytes + len(compressed) <= _GZIP_CACHE_MAX_BYTES: _gzip_cache[cache_key] = compressed _gzip_cache_bytes += len(compressed) return _apply_gzip(response, compressed) def _apply_gzip(response, compressed): response.set_data(compressed) response.headers['Content-Encoding'] = 'gzip' response.headers['Content-Length'] = str(len(compressed)) if 'accept-encoding' not in (response.headers.get('Vary') or '').lower(): response.headers.add('Vary', 'Accept-Encoding') etag = response.headers.get('ETag') if etag and 'gzip' not in etag: response.headers['ETag'] = etag.rstrip('"') + '-gzip"' return response # Add security headers and caching to all responses @app.after_request def add_security_headers(response): """Add security headers and caching to all responses""" # Only set standard security headers - avoid Permissions-Policy to prevent browser warnings # about unrecognized features response.headers['X-Content-Type-Options'] = 'nosniff' response.headers['X-Frame-Options'] = 'SAMEORIGIN' response.headers['X-XSS-Protection'] = '1; mode=block' # Add caching headers for static assets if request.path.startswith(_VERSIONED_ASSET_PREFIXES): # Cache static assets for 1 year (with versioning via query params) response.headers['Cache-Control'] = 'public, max-age=31536000, immutable' response.headers['Expires'] = (datetime.now() + timedelta(days=365)).strftime('%a, %d %b %Y %H:%M:%S GMT') elif request.path.startswith('/api/v3/'): if request.method == 'GET' and 'stream' not in request.path: if response.mimetype == 'application/json': # JSON is live state. A cached copy made a fetch() right after # an install, toggle or Wi-Fi connect show the state from # before it for up to the cache lifetime. response.headers['Cache-Control'] = 'no-store' else: # Files served through the API (plugin web UI assets, backup # downloads) keep the short cache. response.headers['Cache-Control'] = 'private, max-age=5, must-revalidate' else: # No cache for HTML pages to ensure fresh content response.headers['Cache-Control'] = 'no-cache, no-store, must-revalidate' response.headers['Pragma'] = 'no-cache' response.headers['Expires'] = '0' return response class _StreamBroadcaster: """Fan-out broadcaster: one background generator thread pushes to all SSE clients. This means N browser tabs share one generator instead of each running their own, keeping PIL encodes / subprocess forks constant regardless of how many tabs are open. """ def __init__(self, generator_factory): self._generator_factory = generator_factory self._clients: set = set() self._lock = threading.Lock() self._thread: threading.Thread | None = None def subscribe(self) -> queue.Queue: q: queue.Queue = queue.Queue(maxsize=5) with self._lock: self._clients.add(q) if not (self._thread and self._thread.is_alive()): self._thread = threading.Thread(target=self._broadcast, daemon=True) self._thread.start() return q def unsubscribe(self, q: queue.Queue) -> None: with self._lock: self._clients.discard(q) def _broadcast(self): for data in self._generator_factory(): with self._lock: if not self._clients: # No subscribers — exit so the thread doesn't spin indefinitely. # subscribe() will restart it when a new client arrives. # Drop the handle here, under the lock: between this break # and the thread actually ending (closing the generator # can take a while) is_alive() is still True, and a # client subscribing then got no thread at all. self._thread = None break for q in self._clients: try: q.put_nowait(data) except queue.Full: # Client is reading too slowly; drop the oldest item and # deliver the latest so the queue never stalls the client. try: q.get_nowait() except queue.Empty: pass try: q.put_nowait(data) except queue.Full: pass def _get_power_status(): """Check Raspberry Pi under-voltage/throttling status via vcgencmd. Returns a dict of decoded flags, or None on non-Pi platforms (no vcgencmd) or if the call fails for any reason. See: https://www.raspberrypi.com/documentation/computers/os.html#get_throttled """ if not _VCGENCMD: return None try: result = subprocess.run( [_VCGENCMD, 'get_throttled'], capture_output=True, text=True, timeout=2 ) match = re.search(r'0x([0-9a-fA-F]+)', result.stdout) if not match: return None bits = int(match.group(1), 16) return { 'under_voltage_now': bool(bits & 0x1), 'freq_capped_now': bool(bits & 0x2), 'throttled_now': bool(bits & 0x4), 'soft_temp_limit_now': bool(bits & 0x8), 'under_voltage_occurred': bool(bits & 0x10000), 'freq_capped_occurred': bool(bits & 0x20000), 'throttled_occurred': bool(bits & 0x40000), 'soft_temp_limit_occurred': bool(bits & 0x80000), } except (subprocess.SubprocessError, OSError, ValueError) as e: app.logger.warning("vcgencmd get_throttled failed: %s", e) return None # System status generator for SSE def system_status_generator(): """Generate system status updates""" while True: try: metrics = collect_system_metrics() cpu_percent = metrics['cpu_percent'] memory_used_percent = metrics['memory_used_percent'] memory_available_mb = metrics['memory_available_mb'] disk_used_percent = metrics['disk_used_percent'] cpu_temp = metrics['cpu_temp'] # Check if display service is running (cached to avoid per-client subprocess forks) service_active = _unit_is_active('ledmatrix', _LEDMATRIX_SERVICE_CACHE_TTL) status = { 'timestamp': time.time(), 'uptime': 'Running', 'service_active': service_active, 'cpu_percent': cpu_percent, 'memory_used_percent': memory_used_percent, 'memory_available_mb': memory_available_mb, 'cpu_temp': cpu_temp, 'disk_used_percent': disk_used_percent, 'power': _get_power_status() } yield status except Exception: app.logger.error("SSE generator error", exc_info=True) yield {'error': 'An error occurred; see server logs'} time.sleep(10) # Update every 10 seconds (reduced frequency for better performance) # Display preview generator for SSE def display_preview_generator(): """Generate display preview updates from snapshot file""" snapshot_path = display_preview.SNAPSHOT_PATH # Viewer marker: this generator only runs while the broadcaster has # subscribers (it exits with no clients), so touching the marker each # loop tells the DISPLAY service a browser is actually watching — it # only pays for full-rate PNG snapshot encodes while this stays fresh # (see src/common/snapshot_policy.py). viewer_marker_path = "/tmp/led_matrix_preview_viewer" # nosec B108 - fixed path matches display_manager last_modified = None def _touch_viewer_marker(): try: with open(viewer_marker_path, 'a'): pass os.utime(viewer_marker_path, None) except OSError: pass # display side treats a missing marker as "no viewer" # Get display dimensions from config: the logical size DisplayManager # renders at, so double-sided setups preview one screen from src.display_geometry import logical_size try: width, height = logical_size(config_manager.load_config()) except (KeyError, TypeError, ValueError, AttributeError, ConfigError): width, height = logical_size({}) while True: try: _touch_viewer_marker() # Check if snapshot file exists and has been modified if os.path.exists(snapshot_path): current_modified = os.path.getmtime(snapshot_path) # Only read if file is new or has been updated if last_modified is None or current_modified > last_modified: try: preview_data = display_preview.preview_payload( width, height, display_preview.read_snapshot_base64(snapshot_path)) last_modified = current_modified yield preview_data except OSError: # Transient filesystem race (file rotated/replaced # between mtime check and read); skip this update. app.logger.debug("Preview snapshot read failed; skipping frame", exc_info=True) else: yield display_preview.preview_payload(width, height, None) except Exception: app.logger.error("SSE generator error", exc_info=True) yield {'error': 'An error occurred; see server logs'} time.sleep(1.0) # the snapshot is re-read only when its mtime changes # Logs generator for SSE def logs_generator(): """Generate log updates from journalctl""" while True: try: # Reading the journal without sudo needs the systemd-journal group. try: if not _JOURNALCTL: yield {'timestamp': time.time(), 'logs': 'journalctl not found; cannot read logs'} time.sleep(60) continue result = subprocess.run( [_JOURNALCTL, '-u', 'ledmatrix.service', '-u', 'ledmatrix-web.service', '-n', '50', '--no-pager', '--output=short-iso'], capture_output=True, text=True, timeout=5 ) if result.returncode == 0: logs_text = result.stdout.strip() if logs_text: logs_data = { 'timestamp': time.time(), 'logs': logs_text } yield logs_data else: # No logs available logs_data = { 'timestamp': time.time(), 'logs': 'No logs available from ledmatrix or ledmatrix-web service' } yield logs_data else: # journalctl failed error_data = { 'timestamp': time.time(), 'logs': f'journalctl failed with return code {result.returncode}: {result.stderr.strip()}' } yield error_data except subprocess.TimeoutExpired: # Timeout - just skip this update pass except Exception: app.logger.error("Error running journalctl", exc_info=True) error_data = { 'timestamp': time.time(), 'logs': 'Error running journalctl; see server logs' } yield error_data except Exception: app.logger.error("Unexpected error in logs generator", exc_info=True) error_data = { 'timestamp': time.time(), 'logs': 'Unexpected error in logs generator; see server logs' } yield error_data time.sleep(5) # Update every 5 seconds (reduced frequency for better performance) # One broadcaster per stream — shared across all SSE clients _stats_broadcaster = _StreamBroadcaster(system_status_generator) _display_broadcaster = _StreamBroadcaster(display_preview_generator) _logs_broadcaster = _StreamBroadcaster(logs_generator) def _sse_stream(broadcaster: _StreamBroadcaster) -> Response: """Return a streaming SSE response backed by a shared broadcaster.""" q = broadcaster.subscribe() def generate(): try: while True: try: data = q.get(timeout=30) yield f"data: {json.dumps(data)}\n\n" except queue.Empty: # Send an SSE comment heartbeat to keep the connection alive # through proxies that close idle connections. yield ": heartbeat\n\n" except GeneratorExit: pass finally: broadcaster.unsubscribe(q) return Response(generate(), mimetype='text/event-stream') # SSE endpoints @app.route('/api/v3/stream/stats') def stream_stats(): return _sse_stream(_stats_broadcaster) @app.route('/api/v3/stream/display') def stream_display(): return _sse_stream(_display_broadcaster) @app.route('/api/v3/stream/logs') def stream_logs(): return _sse_stream(_logs_broadcaster) # Each SSE stream is one long-lived request, so only a (re)connect counts # against a limit. The streams get their own 200 per minute, tighter than the # 1000 per minute default, which bounds a client stuck reconnecting. # flask-limiter enforces a decorated limit in the wrapper limit() returns, and # marks the original function exempt from the default, so the wrapper has to # replace the registered view: discarding it leaves the streams unlimited. if limiter: for _endpoint in ('stream_stats', 'stream_display', 'stream_logs'): app.view_functions[_endpoint] = limiter.limit("200 per minute")( app.view_functions[_endpoint] ) @app.route('/favicon.ico') def favicon(): """Return 204 No Content for favicon to avoid 404 errors""" return '', 204 _reconciliation_started = False _reconciliation_lock = threading.Lock() def _run_startup_reconciliation() -> None: """Run state reconciliation in background to auto-repair missing plugins. Runs once per process, whether or not every inconsistency could be fixed: ``_reconciliation_started`` is never reset. Resetting it after a failed repair (a config entry naming a plugin the registry no longer has) made the ``before_request`` hook rerun reconciliation on every request, an install-retry loop that pegged the CPU and flooded the log. Unresolved issues are left for the user to address in the UI. """ from src.logging_config import get_logger _logger = get_logger('reconciliation') try: from src.plugin_system.state_reconciliation import StateReconciliation from src.plugin_system.plugin_runtime import read_plugin_runtime reconciler = StateReconciliation( config_manager=config_manager, plugins_dir=plugins_dir, store_manager=plugin_store_manager, runtime_source=lambda: read_plugin_runtime(api_v3.cache_manager), ) result = reconciler.reconcile_state() if result.inconsistencies_found: _logger.info("[Reconciliation] %s", result.message) if result.inconsistencies_fixed: plugin_catalog.discover_plugins() if not result.reconciliation_successful: _logger.warning( "[Reconciliation] Finished with %d unresolved issue(s); " "will not retry automatically. Use the Plugin Store or the " "manual 'Reconcile' action to resolve.", len(result.inconsistencies_manual), ) # Write status file so the web UI can surface unresolved issues as a # banner without the user having to read journalctl. Mirrors the # hw_status pattern (/tmp/led_matrix_hw_status.json). import tempfile _recon_status = { "done": True, "successful": result.reconciliation_successful, "fixed_count": len(result.inconsistencies_fixed), "unresolved": [ { "plugin_id": inc.plugin_id, "type": inc.inconsistency_type.value, "description": inc.description, } for inc in result.inconsistencies_manual ], } _recon_path = os.path.join(tempfile.gettempdir(), "ledmatrix_reconciliation.json") _tmp = None try: if not os.path.islink(_recon_path): _fd, _tmp = tempfile.mkstemp(dir=tempfile.gettempdir(), prefix=".led_recon_") with os.fdopen(_fd, "w") as _f: json.dump(_recon_status, _f) os.replace(_tmp, _recon_path) _tmp = None # Rename succeeded; nothing to clean up except (OSError, ValueError, TypeError) as _e: _logger.warning("[Reconciliation] Could not write status file: %s", _e) finally: if _tmp is not None and os.path.exists(_tmp): try: os.unlink(_tmp) except OSError: pass except Exception as e: _logger.error("[Reconciliation] Error: %s", e, exc_info=True) # Run reconciliation in the background on first request @app.before_request def start_startup_reconciliation(): """Launch startup reconciliation in the background once.""" global _reconciliation_started with _reconciliation_lock: if not _reconciliation_started: _reconciliation_started = True threading.Thread(target=_run_startup_reconciliation, daemon=True).start() _auto_updater = None def start_auto_update_scheduler(): """Start the weekly auto-update thread (web_interface/auto_update.py). Called by the launchers, not at import: tests and dev tools import this module, and none of them should ever git pull. Idle unless auto_update.enabled is set, so it is safe to always start. """ global _auto_updater if _auto_updater is not None: return _auto_updater from web_interface.auto_update import AutoUpdater from web_interface.blueprints.api_v3.system import perform_core_update _auto_updater = AutoUpdater( config_manager=config_manager, core_update=perform_core_update, store_manager=plugin_store_manager, plugin_catalog=plugin_catalog, schema_manager=schema_manager, operation_history=operation_history, ) _auto_updater.start() return _auto_updater if __name__ == '__main__': start_auto_update_scheduler() # threaded=True is Flask's default since 1.0 but stated explicitly so that # long-lived /api/v3/stream/* SSE connections don't starve other requests. # Debug mode is off by default; opt in with FLASK_DEBUG=1 in the environment. _debug = os.environ.get('FLASK_DEBUG', '0') == '1' app.run(host='0.0.0.0', port=5000, debug=_debug, threaded=True) # nosec B104 - intentional; local network device