mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-04 06:15:09 +00:00
fix(display): thread-safety for deferred updates, BDF faces and follower image; one refresh default (#652)
- DisplayManager.defer_update()/process_deferred_updates(): one lock around every queue mutation (appends from the update thread were lost to the render thread's filter/slice reassignments); callables run outside it. - FontManager and element_style no longer cache BDF freetype.Face objects process-wide (load_bdf_face caches them per thread); element_style's LRU is locked against get/move_to_end vs eviction races. - limit_refresh_rate_hz default is one constant, DEFAULT_REFRESH_LIMIT_HZ = 100 (the template's), for the library options, refresh_hz, the matrix guard, Vegas and scroll_config. Previously a missing key capped the panel at 90 while pacing assumed 100. - Sync follower: the TCP thread queues the leader's scroll image; the render thread swaps image/array/width in between frames. - update_display() error log rate-limited (traceback first, then once a minute with a count); swallowed DisplayController exceptions log at DEBUG. - Root display_controller.py runs run.py via runpy. - stream_manager: correct the RLock release comments; merge duplicate if. Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -19,6 +19,13 @@ accepts both, but the store flags the old spelling as deprecated
|
|||||||
|
|
||||||
## Unreleased
|
## Unreleased
|
||||||
|
|
||||||
|
- Display thread-safety and consistency fixes:
|
||||||
|
- `DisplayManager.defer_update()` from a plugin's update thread no longer loses queued updates while the render thread processes the queue; the queue is locked, and the queued callables still run outside the lock.
|
||||||
|
- BDF fonts: `FontManager.get_font()` and `element_style.load_font()` no longer hand one `freetype.Face` to every thread. BDF faces come from `load_bdf_face`, which already caches them per thread; TrueType fonts are cached as before. `element_style`'s font cache is locked (a concurrent eviction could raise `KeyError`).
|
||||||
|
- **Behaviour change:** when `display.hardware.limit_refresh_rate_hz` is missing from config, the panel is now capped at 100 Hz (the config template's value) instead of 90 Hz. Scroll pacing already assumed 100 Hz in that case, so it now matches what the panel does. Configs that set the key (every config migrated from the template) are unaffected.
|
||||||
|
- A sync follower adopts the leader's scroll image between frames on the render thread, instead of the TCP thread swapping the image, array and width while a frame is being drawn.
|
||||||
|
- `update_display()` errors are logged once with a traceback, then at most once a minute with a count, instead of an untraced line every frame. Several swallowed exceptions in `DisplayController` now log at DEBUG.
|
||||||
|
- The repo-root `display_controller.py` now runs `run.py` (the real entry point), so it gets run.py's `-e`/`-d` flags, logging setup and `sys.dont_write_bytecode`.
|
||||||
- Vegas: a plugin set to `vegas_mode: "static"` pauses the scroll for its turn again. The pause was triggered by peeking at the front of a segment buffer that continuous scrolling (the default) never advances, so a static plugin paused only if it happened to be first, once, at startup, and otherwise scrolled past as ordinary content; swap mode had the same problem for any static plugin not first in its cycle. The render pipeline now marks where each static plugin's turn falls in the strip and the scroll pauses when it gets there. The pause runs the plugin's `display()` under its plugin lock, and a static plugin's content is no longer rendered for the strip.
|
- Vegas: a plugin set to `vegas_mode: "static"` pauses the scroll for its turn again. The pause was triggered by peeking at the front of a segment buffer that continuous scrolling (the default) never advances, so a static plugin paused only if it happened to be first, once, at startup, and otherwise scrolled past as ordinary content; swap mode had the same problem for any static plugin not first in its cycle. The render pipeline now marks where each static plugin's turn falls in the strip and the scroll pauses when it gets there. The pause runs the plugin's `display()` under its plugin lock, and a static plugin's content is no longer rendered for the strip.
|
||||||
|
|
||||||
- The display loop no longer spins at 100% CPU when no enabled mode has anything to show (for example, only a sports plugin enabled in its off-season). After one full rotation of empty modes it checks one mode per second until something shows; live content still takes over at once.
|
- The display loop no longer spins at 100% CPU when no enabled mode has anything to show (for example, only a sports plugin enabled in its off-season). After one full rotation of empty modes it checks one mode per second until something shows; live content still takes over at once.
|
||||||
|
|||||||
+15
-7
@@ -1,12 +1,20 @@
|
|||||||
#!/usr/bin/env python3
|
#!/usr/bin/env python3
|
||||||
|
"""Legacy entry point: runs ``run.py``, which is the one to use.
|
||||||
|
|
||||||
|
``python3 run.py`` (``-e`` for the emulator, ``-d`` for debug logging) is how
|
||||||
|
the display service and the docs start LEDMatrix. This file used to import
|
||||||
|
``src.display_controller.main`` directly, which skipped what run.py sets up
|
||||||
|
first -- ``sys.dont_write_bytecode`` (root-owned ``__pycache__`` in plugin
|
||||||
|
directories blocks the web service from updating them), the ``-e``/``-d``
|
||||||
|
flags, and the logging configuration. It now runs run.py exactly as
|
||||||
|
``python3 run.py`` would, with the same arguments.
|
||||||
|
"""
|
||||||
|
|
||||||
import os
|
import os
|
||||||
import sys
|
import runpy
|
||||||
|
|
||||||
# Add the project root directory to Python path
|
|
||||||
sys.path.append(os.path.dirname(os.path.abspath(__file__)))
|
|
||||||
|
|
||||||
from src.display_controller import main
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
main()
|
runpy.run_path(
|
||||||
|
os.path.join(os.path.dirname(os.path.abspath(__file__)), "run.py"),
|
||||||
|
run_name="__main__",
|
||||||
|
)
|
||||||
|
|||||||
@@ -51,6 +51,8 @@ import logging
|
|||||||
from dataclasses import dataclass, replace
|
from dataclasses import dataclass, replace
|
||||||
from typing import Any, Dict, Optional
|
from typing import Any, Dict, Optional
|
||||||
|
|
||||||
|
from src.matrix_support import DEFAULT_REFRESH_LIMIT_HZ
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
#: Speed used when a plugin supplies nothing usable. One pixel per refresh on a
|
#: Speed used when a plugin supplies nothing usable. One pixel per refresh on a
|
||||||
@@ -63,8 +65,9 @@ MIN_PIXELS_PER_SECOND = 1.0
|
|||||||
MAX_PIXELS_PER_SECOND = 500.0
|
MAX_PIXELS_PER_SECOND = 500.0
|
||||||
|
|
||||||
#: Assumed refresh when the caller does not say. Matches the usual
|
#: Assumed refresh when the caller does not say. Matches the usual
|
||||||
#: ``display.hardware.limit_refresh_rate_hz``.
|
#: ``display.hardware.limit_refresh_rate_hz``, and is the cap DisplayManager
|
||||||
DEFAULT_REFRESH_HZ = 100.0
|
#: applies when that key is missing.
|
||||||
|
DEFAULT_REFRESH_HZ = float(DEFAULT_REFRESH_LIMIT_HZ)
|
||||||
|
|
||||||
#: How far px/s may sit from a whole number of pixels per refresh before it is
|
#: How far px/s may sit from a whole number of pixels per refresh before it is
|
||||||
#: worth warning about. 0.05px per frame is invisible; a third of a pixel is not.
|
#: worth warning about. 0.05px per frame is invisible; a third of a pixel is not.
|
||||||
|
|||||||
+46
-14
@@ -27,6 +27,7 @@ import signal
|
|||||||
import json
|
import json
|
||||||
import threading
|
import threading
|
||||||
import types
|
import types
|
||||||
|
from collections import deque
|
||||||
from contextlib import contextmanager
|
from contextlib import contextmanager
|
||||||
from typing import Dict, Any, List, Optional, Callable, Tuple
|
from typing import Dict, Any, List, Optional, Callable, Tuple
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
@@ -196,6 +197,12 @@ class DisplayController:
|
|||||||
# scroll image arrives.
|
# scroll image arrives.
|
||||||
self._follower_pending_new_image = False
|
self._follower_pending_new_image = False
|
||||||
self._follower_last_frame = None
|
self._follower_last_frame = None
|
||||||
|
# (image, array) from the leader, handed from the sync TCP thread to
|
||||||
|
# the render thread, which adopts it at the start of a follower frame
|
||||||
|
# (_adopt_follower_scroll_image). One append / one popleft, each
|
||||||
|
# atomic, so the render thread never draws from a half-swapped
|
||||||
|
# cached_image / cached_array / total_scroll_width.
|
||||||
|
self._follower_incoming_image: deque = deque(maxlen=1)
|
||||||
self._follower_deadline: Optional[float] = None
|
self._follower_deadline: Optional[float] = None
|
||||||
# Leader: time.time() of the last follower frame sent.
|
# Leader: time.time() of the last follower frame sent.
|
||||||
self._last_follower_send = 0.0
|
self._last_follower_send = 0.0
|
||||||
@@ -559,21 +566,15 @@ class DisplayController:
|
|||||||
# temporarily replace the leader's correct one.
|
# temporarily replace the leader's correct one.
|
||||||
|
|
||||||
# When the leader sends its scroll image (TCP), update our
|
# When the leader sends its scroll image (TCP), update our
|
||||||
# cached_array so both Pis have pixel-identical images.
|
# cached_array so both Pis have pixel-identical images. This runs
|
||||||
|
# on the sync TCP thread, so it only converts and queues the
|
||||||
|
# image; the render thread swaps it in between frames.
|
||||||
import numpy as _np
|
import numpy as _np
|
||||||
def _on_leader_scroll_image(image):
|
def _on_leader_scroll_image(image):
|
||||||
vc = self.vegas_coordinator
|
vc = self.vegas_coordinator
|
||||||
if vc and vc.render_pipeline:
|
if vc and vc.render_pipeline:
|
||||||
rp = vc.render_pipeline
|
|
||||||
arr = _np.asarray(image.convert("RGB"), dtype=_np.uint8)
|
arr = _np.asarray(image.convert("RGB"), dtype=_np.uint8)
|
||||||
rp.scroll_helper.cached_image = image
|
self._follower_incoming_image.append((image, arr))
|
||||||
rp.scroll_helper.cached_array = arr
|
|
||||||
rp.scroll_helper.total_scroll_width = image.width
|
|
||||||
self._follower_pending_new_image = False
|
|
||||||
logger.info(
|
|
||||||
"Sync: follower adopted leader scroll image %dx%d",
|
|
||||||
image.width, image.height,
|
|
||||||
)
|
|
||||||
self.sync_manager.set_on_scroll_image(_on_leader_scroll_image)
|
self.sync_manager.set_on_scroll_image(_on_leader_scroll_image)
|
||||||
|
|
||||||
if self.sync_manager.role == SyncRole.LEADER:
|
if self.sync_manager.role == SyncRole.LEADER:
|
||||||
@@ -596,6 +597,29 @@ class DisplayController:
|
|||||||
logger.error("Failed to initialize Vegas mode: %s", e, exc_info=True)
|
logger.error("Failed to initialize Vegas mode: %s", e, exc_info=True)
|
||||||
self.vegas_coordinator = None
|
self.vegas_coordinator = None
|
||||||
|
|
||||||
|
def _adopt_follower_scroll_image(self, rp) -> None:
|
||||||
|
"""Swap in the leader's latest scroll image, on the render thread.
|
||||||
|
|
||||||
|
The sync TCP thread used to set cached_image, cached_array and
|
||||||
|
total_scroll_width one after another while this thread read them, so
|
||||||
|
a frame could slice the new array with the old width. It now queues
|
||||||
|
the image and this applies it between frames.
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
image, arr = self._follower_incoming_image.popleft()
|
||||||
|
except IndexError:
|
||||||
|
return
|
||||||
|
if rp is None:
|
||||||
|
return
|
||||||
|
rp.scroll_helper.cached_image = image
|
||||||
|
rp.scroll_helper.cached_array = arr
|
||||||
|
rp.scroll_helper.total_scroll_width = image.width
|
||||||
|
self._follower_pending_new_image = False
|
||||||
|
logger.info(
|
||||||
|
"Sync: follower adopted leader scroll image %dx%d",
|
||||||
|
image.width, image.height,
|
||||||
|
)
|
||||||
|
|
||||||
def _is_vegas_mode_active(self) -> bool:
|
def _is_vegas_mode_active(self) -> bool:
|
||||||
"""Check if Vegas mode should be running."""
|
"""Check if Vegas mode should be running."""
|
||||||
self._apply_pending_vegas_init()
|
self._apply_pending_vegas_init()
|
||||||
@@ -1667,7 +1691,10 @@ class DisplayController:
|
|||||||
if plugin_instance.has_live_content():
|
if plugin_instance.has_live_content():
|
||||||
live_with_content.append(live_mode)
|
live_with_content.append(live_mode)
|
||||||
except Exception:
|
except Exception:
|
||||||
pass
|
# Treated as no live content; logged so a plugin whose
|
||||||
|
# check always raises is findable.
|
||||||
|
logger.debug("has_live_content() failed for %s", live_mode,
|
||||||
|
exc_info=True)
|
||||||
|
|
||||||
# Build mode list: live modes with content first, then other modes, then live modes without content
|
# Build mode list: live modes with content first, then other modes, then live modes without content
|
||||||
if live_with_content:
|
if live_with_content:
|
||||||
@@ -1904,7 +1931,9 @@ class DisplayController:
|
|||||||
if bg_service and hasattr(bg_service, 'log_memory_stats'):
|
if bg_service and hasattr(bg_service, 'log_memory_stats'):
|
||||||
bg_service.log_memory_stats()
|
bg_service.log_memory_stats()
|
||||||
except Exception:
|
except Exception:
|
||||||
pass # Background service may not be initialized
|
# Background service may not be initialized
|
||||||
|
logger.debug("Background service memory stats unavailable",
|
||||||
|
exc_info=True)
|
||||||
|
|
||||||
# Log deferred updates stats
|
# Log deferred updates stats
|
||||||
if hasattr(self.display_manager, '_scrolling_state'):
|
if hasattr(self.display_manager, '_scrolling_state'):
|
||||||
@@ -2106,6 +2135,7 @@ class DisplayController:
|
|||||||
vc = self.vegas_coordinator
|
vc = self.vegas_coordinator
|
||||||
rp = vc.render_pipeline if (vc and vc.render_pipeline) else None
|
rp = vc.render_pipeline if (vc and vc.render_pipeline) else None
|
||||||
width = self.display_manager.width
|
width = self.display_manager.width
|
||||||
|
self._adopt_follower_scroll_image(rp)
|
||||||
|
|
||||||
local_x = self._follower_local_x
|
local_x = self._follower_local_x
|
||||||
if local_x is None:
|
if local_x is None:
|
||||||
@@ -2876,7 +2906,8 @@ class DisplayController:
|
|||||||
try:
|
try:
|
||||||
self.wifi_status_file.unlink()
|
self.wifi_status_file.unlink()
|
||||||
except Exception:
|
except Exception:
|
||||||
pass
|
logger.debug("Could not remove WiFi status file %s",
|
||||||
|
self.wifi_status_file, exc_info=True)
|
||||||
return None
|
return None
|
||||||
|
|
||||||
# Validate required fields
|
# Validate required fields
|
||||||
@@ -2909,7 +2940,8 @@ class DisplayController:
|
|||||||
try:
|
try:
|
||||||
self.wifi_status_file.unlink()
|
self.wifi_status_file.unlink()
|
||||||
except Exception:
|
except Exception:
|
||||||
pass
|
logger.debug("Could not remove WiFi status file %s",
|
||||||
|
self.wifi_status_file, exc_info=True)
|
||||||
return None
|
return None
|
||||||
|
|
||||||
# Message is valid and not expired — cache for the throttle window
|
# Message is valid and not expired — cache for the throttle window
|
||||||
|
|||||||
+74
-30
@@ -44,7 +44,9 @@ from src.display_geometry import (
|
|||||||
DEFAULT_CHAIN_LENGTH, DEFAULT_COLS, DEFAULT_PARALLEL, DEFAULT_ROWS,
|
DEFAULT_CHAIN_LENGTH, DEFAULT_COLS, DEFAULT_PARALLEL, DEFAULT_ROWS,
|
||||||
compose_pixel_mapper_config, physical_size, resolve_double_sided,
|
compose_pixel_mapper_config, physical_size, resolve_double_sided,
|
||||||
)
|
)
|
||||||
from src.matrix_support import MatrixSettingsRefused, library_refusals, refusal_message
|
from src.matrix_support import (
|
||||||
|
DEFAULT_REFRESH_LIMIT_HZ, MatrixSettingsRefused, library_refusals, refusal_message,
|
||||||
|
)
|
||||||
from src.pi5_matrix_support import is_raspberry_pi_5
|
from src.pi5_matrix_support import is_raspberry_pi_5
|
||||||
import threading
|
import threading
|
||||||
import time
|
import time
|
||||||
@@ -72,6 +74,10 @@ logger = get_logger(__name__)
|
|||||||
#: and therefore get_font_height() -- report anything but 0.
|
#: and therefore get_font_height() -- report anything but 0.
|
||||||
_CALENDAR_FONT_PX = 7
|
_CALENDAR_FONT_PX = 7
|
||||||
|
|
||||||
|
#: Seconds between repeats of update_display()'s error log. It runs every
|
||||||
|
#: frame, so a fault that persists would otherwise log ~100 lines a second.
|
||||||
|
_UPDATE_ERROR_LOG_INTERVAL = 60.0
|
||||||
|
|
||||||
|
|
||||||
def _bdf_native_size(face) -> int:
|
def _bdf_native_size(face) -> int:
|
||||||
"""The pixel height a BDF Face declares, or 0 if it does not say.
|
"""The pixel height a BDF Face declares, or 0 if it does not say.
|
||||||
@@ -236,6 +242,11 @@ class DisplayManager:
|
|||||||
|
|
||||||
_instance = None
|
_instance = None
|
||||||
|
|
||||||
|
# update_display()'s error-log throttle. Class defaults so instances built
|
||||||
|
# without __init__ (tests, doubles) have them too.
|
||||||
|
_update_error_logged_at: Optional[float] = None
|
||||||
|
_update_errors_suppressed = 0
|
||||||
|
|
||||||
def __new__(cls, *args, **kwargs):
|
def __new__(cls, *args, **kwargs):
|
||||||
if cls._instance is None:
|
if cls._instance is None:
|
||||||
cls._instance = super(DisplayManager, cls).__new__(cls)
|
cls._instance = super(DisplayManager, cls).__new__(cls)
|
||||||
@@ -345,6 +356,13 @@ class DisplayManager:
|
|||||||
'max_deferred_updates': 50, # Limit queue size to prevent memory issues
|
'max_deferred_updates': 50, # Limit queue size to prevent memory issues
|
||||||
'deferred_update_ttl': 300.0 # 5 minutes TTL for deferred updates
|
'deferred_update_ttl': 300.0 # 5 minutes TTL for deferred updates
|
||||||
}
|
}
|
||||||
|
# Guards _scrolling_state['deferred_updates']. defer_update() is called
|
||||||
|
# from plugin update() on the update worker thread while
|
||||||
|
# process_deferred_updates() runs on the render thread, and both
|
||||||
|
# rebuild the list (TTL filter, [n:] slice) and assign it back -- an
|
||||||
|
# append landing between one side's read and its assignment was lost.
|
||||||
|
# Never held while a queued callable runs: those may defer again.
|
||||||
|
self._deferred_lock = threading.Lock()
|
||||||
|
|
||||||
self._setup_matrix()
|
self._setup_matrix()
|
||||||
logger.info("Matrix setup completed in %.3f seconds", time.time() - start_time)
|
logger.info("Matrix setup completed in %.3f seconds", time.time() - start_time)
|
||||||
@@ -976,7 +994,21 @@ class DisplayManager:
|
|||||||
# Write a snapshot for the web preview (throttled)
|
# Write a snapshot for the web preview (throttled)
|
||||||
self._write_snapshot_if_due(frame_checksum)
|
self._write_snapshot_if_due(frame_checksum)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"Error updating display: {e}")
|
# Once with the traceback, then at most every
|
||||||
|
# _UPDATE_ERROR_LOG_INTERVAL with a count of what was skipped.
|
||||||
|
now = time.monotonic()
|
||||||
|
last = self._update_error_logged_at
|
||||||
|
if last is None:
|
||||||
|
self._update_error_logged_at = now
|
||||||
|
logger.error("Error updating display: %s", e, exc_info=True)
|
||||||
|
elif now - last >= _UPDATE_ERROR_LOG_INTERVAL:
|
||||||
|
skipped = self._update_errors_suppressed
|
||||||
|
self._update_error_logged_at = now
|
||||||
|
self._update_errors_suppressed = 0
|
||||||
|
logger.error("Error updating display: %s (%d more since the "
|
||||||
|
"last report)", e, skipped)
|
||||||
|
else:
|
||||||
|
self._update_errors_suppressed += 1
|
||||||
|
|
||||||
def _setup_scan_order_compensation(self) -> None:
|
def _setup_scan_order_compensation(self) -> None:
|
||||||
"""Work out which rows to show a refresh behind while scrolling.
|
"""Work out which rows to show a refresh behind while scrolling.
|
||||||
@@ -1547,7 +1579,7 @@ class DisplayManager:
|
|||||||
options.panel_type = hardware_config.get('panel_type', '')
|
options.panel_type = hardware_config.get('panel_type', '')
|
||||||
options.disable_hardware_pulsing = hardware_config.get('disable_hardware_pulsing', False)
|
options.disable_hardware_pulsing = hardware_config.get('disable_hardware_pulsing', False)
|
||||||
options.show_refresh_rate = hardware_config.get('show_refresh_rate', False)
|
options.show_refresh_rate = hardware_config.get('show_refresh_rate', False)
|
||||||
options.limit_refresh_rate_hz = hardware_config.get('limit_refresh_rate_hz', 90)
|
options.limit_refresh_rate_hz = hardware_config.get('limit_refresh_rate_hz', DEFAULT_REFRESH_LIMIT_HZ)
|
||||||
options.gpio_slowdown = runtime_config.get('gpio_slowdown', 3)
|
options.gpio_slowdown = runtime_config.get('gpio_slowdown', 3)
|
||||||
|
|
||||||
# Disable internal privilege dropping - we manage this via systemd or remain root
|
# Disable internal privilege dropping - we manage this via systemd or remain root
|
||||||
@@ -1593,7 +1625,7 @@ class DisplayManager:
|
|||||||
value = float(hardware.get('limit_refresh_rate_hz') or 0)
|
value = float(hardware.get('limit_refresh_rate_hz') or 0)
|
||||||
except (TypeError, ValueError):
|
except (TypeError, ValueError):
|
||||||
value = 0.0
|
value = 0.0
|
||||||
return value if value > 0 else 100.0
|
return value if value > 0 else float(DEFAULT_REFRESH_LIMIT_HZ)
|
||||||
|
|
||||||
def _scrolling_now(self) -> bool:
|
def _scrolling_now(self) -> bool:
|
||||||
"""Whether a scroll is running, without is_currently_scrolling()'s
|
"""Whether a scroll is running, without is_currently_scrolling()'s
|
||||||
@@ -1704,26 +1736,28 @@ class DisplayManager:
|
|||||||
"""
|
"""
|
||||||
current_time = time.time()
|
current_time = time.time()
|
||||||
|
|
||||||
# Clean up expired updates before adding new ones
|
with self._deferred_lock:
|
||||||
self._cleanup_expired_deferred_updates(current_time)
|
# Clean up expired updates before adding new ones
|
||||||
|
self._cleanup_expired_deferred_updates(current_time)
|
||||||
|
|
||||||
# Limit queue size to prevent memory issues
|
# Limit queue size to prevent memory issues
|
||||||
if len(self._scrolling_state['deferred_updates']) >= self._scrolling_state['max_deferred_updates']:
|
if len(self._scrolling_state['deferred_updates']) >= self._scrolling_state['max_deferred_updates']:
|
||||||
# Remove oldest update to make room
|
# Remove oldest update to make room
|
||||||
self._scrolling_state['deferred_updates'].pop(0)
|
self._scrolling_state['deferred_updates'].pop(0)
|
||||||
logger.debug("Removed oldest deferred update due to queue size limit")
|
logger.debug("Removed oldest deferred update due to queue size limit")
|
||||||
|
|
||||||
self._scrolling_state['deferred_updates'].append({
|
self._scrolling_state['deferred_updates'].append({
|
||||||
'func': update_func,
|
'func': update_func,
|
||||||
'priority': priority,
|
'priority': priority,
|
||||||
'timestamp': current_time
|
'timestamp': current_time
|
||||||
})
|
})
|
||||||
|
|
||||||
# Only sort if we have a reasonable number of updates to avoid excessive sorting
|
# Only sort if we have a reasonable number of updates to avoid excessive sorting
|
||||||
if len(self._scrolling_state['deferred_updates']) <= 20:
|
if len(self._scrolling_state['deferred_updates']) <= 20:
|
||||||
self._scrolling_state['deferred_updates'].sort(key=lambda x: x['priority'])
|
self._scrolling_state['deferred_updates'].sort(key=lambda x: x['priority'])
|
||||||
|
|
||||||
logger.debug(f"Deferred update added. Total deferred: {len(self._scrolling_state['deferred_updates'])}")
|
queued = len(self._scrolling_state['deferred_updates'])
|
||||||
|
logger.debug(f"Deferred update added. Total deferred: {queued}")
|
||||||
|
|
||||||
def process_deferred_updates(self):
|
def process_deferred_updates(self):
|
||||||
"""Process any deferred updates if not currently scrolling."""
|
"""Process any deferred updates if not currently scrolling."""
|
||||||
@@ -1731,21 +1765,26 @@ class DisplayManager:
|
|||||||
|
|
||||||
# Always clean up expired updates, even if scrolling
|
# Always clean up expired updates, even if scrolling
|
||||||
# This prevents memory leaks from accumulated expired updates
|
# This prevents memory leaks from accumulated expired updates
|
||||||
self._cleanup_expired_deferred_updates(current_time)
|
with self._deferred_lock:
|
||||||
|
self._cleanup_expired_deferred_updates(current_time)
|
||||||
|
|
||||||
if self.is_currently_scrolling():
|
if self.is_currently_scrolling():
|
||||||
return
|
return
|
||||||
|
|
||||||
if not self._scrolling_state['deferred_updates']:
|
with self._deferred_lock:
|
||||||
return
|
if not self._scrolling_state['deferred_updates']:
|
||||||
|
return
|
||||||
|
|
||||||
# Process only a limited number of updates per call to avoid blocking
|
# Process only a limited number of updates per call to avoid blocking
|
||||||
max_updates_per_call = min(5, len(self._scrolling_state['deferred_updates']))
|
max_updates_per_call = min(5, len(self._scrolling_state['deferred_updates']))
|
||||||
updates_to_process = self._scrolling_state['deferred_updates'][:max_updates_per_call]
|
updates_to_process = self._scrolling_state['deferred_updates'][:max_updates_per_call]
|
||||||
self._scrolling_state['deferred_updates'] = self._scrolling_state['deferred_updates'][max_updates_per_call:]
|
self._scrolling_state['deferred_updates'] = self._scrolling_state['deferred_updates'][max_updates_per_call:]
|
||||||
|
queued = len(self._scrolling_state['deferred_updates'])
|
||||||
|
|
||||||
logger.debug(f"Processing {len(updates_to_process)} deferred updates (queue size: {len(self._scrolling_state['deferred_updates'])})")
|
logger.debug(f"Processing {len(updates_to_process)} deferred updates (queue size: {queued})")
|
||||||
|
|
||||||
|
# The callables run outside the lock: they are plugin code of any
|
||||||
|
# length, and one that defers again would deadlock on it.
|
||||||
failed_updates = []
|
failed_updates = []
|
||||||
for update_info in updates_to_process:
|
for update_info in updates_to_process:
|
||||||
try:
|
try:
|
||||||
@@ -1764,10 +1803,15 @@ class DisplayManager:
|
|||||||
|
|
||||||
# Re-add failed updates to the end of the queue (not the beginning)
|
# Re-add failed updates to the end of the queue (not the beginning)
|
||||||
if failed_updates:
|
if failed_updates:
|
||||||
self._scrolling_state['deferred_updates'].extend(failed_updates)
|
with self._deferred_lock:
|
||||||
|
self._scrolling_state['deferred_updates'].extend(failed_updates)
|
||||||
|
|
||||||
def _cleanup_expired_deferred_updates(self, current_time: float):
|
def _cleanup_expired_deferred_updates(self, current_time: float):
|
||||||
"""Remove expired deferred updates to prevent memory leaks."""
|
"""Remove expired deferred updates to prevent memory leaks.
|
||||||
|
|
||||||
|
Callers hold ``_deferred_lock``: this reads the list and assigns a
|
||||||
|
filtered copy back.
|
||||||
|
"""
|
||||||
ttl = self._scrolling_state['deferred_update_ttl']
|
ttl = self._scrolling_state['deferred_update_ttl']
|
||||||
initial_count = len(self._scrolling_state['deferred_updates'])
|
initial_count = len(self._scrolling_state['deferred_updates'])
|
||||||
|
|
||||||
|
|||||||
+33
-14
@@ -48,6 +48,7 @@ import json
|
|||||||
import logging
|
import logging
|
||||||
import math
|
import math
|
||||||
import os
|
import os
|
||||||
|
import threading
|
||||||
from collections import OrderedDict
|
from collections import OrderedDict
|
||||||
from dataclasses import dataclass
|
from dataclasses import dataclass
|
||||||
from typing import Any, Dict, Optional, Tuple, Union
|
from typing import Any, Dict, Optional, Tuple, Union
|
||||||
@@ -68,8 +69,10 @@ _FONTS_SUBDIR = os.path.join('assets', 'fonts')
|
|||||||
_FALLBACK_FONT_NAME = 'PressStart2P-Regular.ttf'
|
_FALLBACK_FONT_NAME = 'PressStart2P-Regular.ttf'
|
||||||
|
|
||||||
# (resolved absolute path, requested size) -> (font face, realised size).
|
# (resolved absolute path, requested size) -> (font face, realised size).
|
||||||
# BDF faces are stateful in principle, but the core's own FontManager shares
|
# TTF only: a BDF ``freetype.Face`` must never be shared between threads
|
||||||
# faces the same way.
|
# (FreeType does not allow it, and ``load_char`` rewrites the face's glyph
|
||||||
|
# slot), and this cache is process-wide. BDF faces come from
|
||||||
|
# ``load_bdf_face`` every time, which already caches them per thread.
|
||||||
#
|
#
|
||||||
# Bounded LRU rather than the unbounded dict this started as: the display
|
# Bounded LRU rather than the unbounded dict this started as: the display
|
||||||
# process runs for weeks, and every config save can introduce a new
|
# process runs for weeks, and every config save can introduce a new
|
||||||
@@ -78,14 +81,28 @@ _FALLBACK_FONT_NAME = 'PressStart2P-Regular.ttf'
|
|||||||
# every other hot cache (display_manager, font_manager, adaptive_layout).
|
# every other hot cache (display_manager, font_manager, adaptive_layout).
|
||||||
_FONT_CACHE_MAX = 256
|
_FONT_CACHE_MAX = 256
|
||||||
_font_cache: 'OrderedDict[Tuple[str, int], Tuple[Any, int]]' = OrderedDict()
|
_font_cache: 'OrderedDict[Tuple[str, int], Tuple[Any, int]]' = OrderedDict()
|
||||||
|
# load_font is called from the display thread and from plugin update threads.
|
||||||
|
# A get() then move_to_end() pair on an unguarded OrderedDict raises KeyError
|
||||||
|
# when another thread evicts the key in between.
|
||||||
|
_font_cache_lock = threading.Lock()
|
||||||
|
|
||||||
|
|
||||||
|
def _cache_get(key: Tuple[str, int]) -> Optional[Tuple[Any, int]]:
|
||||||
|
"""The cached entry for ``key`` (marked most recently used), or None."""
|
||||||
|
with _font_cache_lock:
|
||||||
|
cached = _font_cache.get(key)
|
||||||
|
if cached is not None:
|
||||||
|
_font_cache.move_to_end(key)
|
||||||
|
return cached
|
||||||
|
|
||||||
|
|
||||||
def _cache_put(key: Tuple[str, int], value: Tuple[Any, int]) -> None:
|
def _cache_put(key: Tuple[str, int], value: Tuple[Any, int]) -> None:
|
||||||
"""Insert, evicting the least recently used entry past the bound."""
|
"""Insert, evicting the least recently used entry past the bound."""
|
||||||
_font_cache[key] = value
|
with _font_cache_lock:
|
||||||
_font_cache.move_to_end(key)
|
_font_cache[key] = value
|
||||||
while len(_font_cache) > _FONT_CACHE_MAX:
|
_font_cache.move_to_end(key)
|
||||||
_font_cache.popitem(last=False)
|
while len(_font_cache) > _FONT_CACHE_MAX:
|
||||||
|
_font_cache.popitem(last=False)
|
||||||
|
|
||||||
# Config keys a style element block carries, in schema/UI order.
|
# Config keys a style element block carries, in schema/UI order.
|
||||||
_STYLE_KEYS = ('font', 'font_size', 'text_color', 'visible', 'align')
|
_STYLE_KEYS = ('font', 'font_size', 'text_color', 'visible', 'align')
|
||||||
@@ -222,14 +239,15 @@ def _load_font_sized(font_name: str, size: int) -> Tuple[Any, int]:
|
|||||||
logger.warning("Font file not found: %s, using fallback", font_name)
|
logger.warning("Font file not found: %s, using fallback", font_name)
|
||||||
return _load_fallback_font(size)
|
return _load_fallback_font(size)
|
||||||
|
|
||||||
|
is_bdf = path.lower().endswith('.bdf')
|
||||||
cache_key = (path, size)
|
cache_key = (path, size)
|
||||||
cached = _font_cache.get(cache_key)
|
if not is_bdf:
|
||||||
if cached is not None:
|
cached = _cache_get(cache_key)
|
||||||
_font_cache.move_to_end(cache_key)
|
if cached is not None:
|
||||||
return cached
|
return cached
|
||||||
|
|
||||||
try:
|
try:
|
||||||
if path.lower().endswith('.bdf'):
|
if is_bdf:
|
||||||
font, effective = _load_bdf(path, size)
|
font, effective = _load_bdf(path, size)
|
||||||
else:
|
else:
|
||||||
font, effective = load_truetype(path, size), size
|
font, effective = load_truetype(path, size), size
|
||||||
@@ -238,7 +256,9 @@ def _load_font_sized(font_name: str, size: int) -> Tuple[Any, int]:
|
|||||||
path, size, e)
|
path, size, e)
|
||||||
return _load_fallback_font(size)
|
return _load_fallback_font(size)
|
||||||
|
|
||||||
_cache_put(cache_key, (font, effective))
|
# Not BDF: load_bdf_face caches those per thread (see _font_cache).
|
||||||
|
if not is_bdf:
|
||||||
|
_cache_put(cache_key, (font, effective))
|
||||||
return font, effective
|
return font, effective
|
||||||
|
|
||||||
|
|
||||||
@@ -247,9 +267,8 @@ def _load_fallback_font(size: int) -> Tuple[Any, int]:
|
|||||||
path = resolve_font_path(_FALLBACK_FONT_NAME)
|
path = resolve_font_path(_FALLBACK_FONT_NAME)
|
||||||
if path is not None:
|
if path is not None:
|
||||||
cache_key = (path, size)
|
cache_key = (path, size)
|
||||||
cached = _font_cache.get(cache_key)
|
cached = _cache_get(cache_key)
|
||||||
if cached is not None:
|
if cached is not None:
|
||||||
_font_cache.move_to_end(cache_key)
|
|
||||||
return cached
|
return cached
|
||||||
try:
|
try:
|
||||||
entry = (load_truetype(path, size), size)
|
entry = (load_truetype(path, size), size)
|
||||||
|
|||||||
+7
-1
@@ -498,6 +498,7 @@ class FontManager:
|
|||||||
self.performance_stats["cache_misses"] += 1
|
self.performance_stats["cache_misses"] += 1
|
||||||
|
|
||||||
# Load font
|
# Load font
|
||||||
|
shareable = True
|
||||||
font_path = self.font_catalog.get(family)
|
font_path = self.font_catalog.get(family)
|
||||||
if not font_path:
|
if not font_path:
|
||||||
logger.warning(f"Font family '{family}' not found")
|
logger.warning(f"Font family '{family}' not found")
|
||||||
@@ -507,6 +508,7 @@ class FontManager:
|
|||||||
try:
|
try:
|
||||||
if font_path.endswith('.bdf'):
|
if font_path.endswith('.bdf'):
|
||||||
font = self._load_bdf_font(font_path, size_px)
|
font = self._load_bdf_font(font_path, size_px)
|
||||||
|
shareable = False
|
||||||
else:
|
else:
|
||||||
font = load_truetype(font_path, size_px)
|
font = load_truetype(font_path, size_px)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
@@ -516,7 +518,11 @@ class FontManager:
|
|||||||
self.performance_stats["failed_loads"] += 1
|
self.performance_stats["failed_loads"] += 1
|
||||||
font = ImageFont.load_default()
|
font = ImageFont.load_default()
|
||||||
|
|
||||||
self.font_cache[cache_key] = font
|
# A BDF face is not cached here: font_cache is shared by every
|
||||||
|
# thread, and a freetype.Face must never be (see load_bdf_face, which
|
||||||
|
# already caches BDF faces per thread).
|
||||||
|
if shareable:
|
||||||
|
self.font_cache[cache_key] = font
|
||||||
return font
|
return font
|
||||||
|
|
||||||
def _load_bdf_font(self, font_path: str, size_px: int) -> freetype.Face:
|
def _load_bdf_font(self, font_path: str, size_px: int) -> freetype.Face:
|
||||||
|
|||||||
@@ -74,6 +74,13 @@ MAPPING_OUTPUTS: Dict[str, int] = {
|
|||||||
'classic-pi1': 1,
|
'classic-pi1': 1,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#: The refresh cap (``display.hardware.limit_refresh_rate_hz``) when config
|
||||||
|
#: omits it -- config/config.template.json's value. DisplayManager passes it to
|
||||||
|
#: the library and reports it as ``refresh_hz`` for scroll pacing, so the two
|
||||||
|
#: must be the same number: they were 90 and 100, and pacing solved against a
|
||||||
|
#: rate the panel was capped below.
|
||||||
|
DEFAULT_REFRESH_LIMIT_HZ = 100
|
||||||
|
|
||||||
#: What DisplayManager passes when a key is missing from display.hardware /
|
#: What DisplayManager passes when a key is missing from display.hardware /
|
||||||
#: display.runtime. Config migration normally fills these from
|
#: display.runtime. Config migration normally fills these from
|
||||||
#: config/config.template.json first, so they rarely apply.
|
#: config/config.template.json first, so they rarely apply.
|
||||||
@@ -81,7 +88,7 @@ DISPLAY_MANAGER_DEFAULTS: Dict[str, Any] = {
|
|||||||
'rows': 32, 'cols': 64, 'chain_length': 2, 'parallel': 1,
|
'rows': 32, 'cols': 64, 'chain_length': 2, 'parallel': 1,
|
||||||
'hardware_mapping': 'adafruit-hat-pwm', 'brightness': 90, 'pwm_bits': 10,
|
'hardware_mapping': 'adafruit-hat-pwm', 'brightness': 90, 'pwm_bits': 10,
|
||||||
'pwm_lsb_nanoseconds': 150, 'led_rgb_sequence': 'RGB',
|
'pwm_lsb_nanoseconds': 150, 'led_rgb_sequence': 'RGB',
|
||||||
'row_address_type': 0, 'multiplexing': 0, 'limit_refresh_rate_hz': 90,
|
'row_address_type': 0, 'multiplexing': 0, 'limit_refresh_rate_hz': DEFAULT_REFRESH_LIMIT_HZ,
|
||||||
'gpio_slowdown': 3,
|
'gpio_slowdown': 3,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -16,6 +16,7 @@ from PIL import Image
|
|||||||
|
|
||||||
from src.common.scroll_config import solve_crisp
|
from src.common.scroll_config import solve_crisp
|
||||||
from src.common.scroll_helper import ScrollHelper
|
from src.common.scroll_helper import ScrollHelper
|
||||||
|
from src.matrix_support import DEFAULT_REFRESH_LIMIT_HZ
|
||||||
from src.vegas_mode.config import VegasModeConfig
|
from src.vegas_mode.config import VegasModeConfig
|
||||||
from src.vegas_mode.geometry import separation_gap
|
from src.vegas_mode.geometry import separation_gap
|
||||||
from src.vegas_mode.stream_manager import StreamManager
|
from src.vegas_mode.stream_manager import StreamManager
|
||||||
@@ -199,7 +200,7 @@ class RenderPipeline:
|
|||||||
hz = float(getattr(self.display_manager, 'refresh_hz', 0) or 0)
|
hz = float(getattr(self.display_manager, 'refresh_hz', 0) or 0)
|
||||||
except (TypeError, ValueError):
|
except (TypeError, ValueError):
|
||||||
hz = 0.0
|
hz = 0.0
|
||||||
return hz if hz > 0 else 100.0
|
return hz if hz > 0 else float(DEFAULT_REFRESH_LIMIT_HZ)
|
||||||
|
|
||||||
def _refresh_hz(self) -> float:
|
def _refresh_hz(self) -> float:
|
||||||
"""The refresh to solve the crisp speed against: measured, else the cap."""
|
"""The refresh to solve the crisp speed against: measured, else the cap."""
|
||||||
|
|||||||
@@ -74,8 +74,11 @@ class StreamManager:
|
|||||||
|
|
||||||
# Segments composed into the current cycle (swap mode only).
|
# Segments composed into the current cycle (swap mode only).
|
||||||
self._active_buffer: Deque[ContentSegment] = deque()
|
self._active_buffer: Deque[ContentSegment] = deque()
|
||||||
# Reentrant: _prefetch_content releases and re-acquires it around the
|
# Reentrant: get_next_segment holds it while calling
|
||||||
# slow fetch while a caller may already hold it.
|
# _prefetch_content, which acquires it again. _prefetch_content's
|
||||||
|
# release() around the slow fetch only frees the lock when its caller
|
||||||
|
# did not already hold it (initialize); from get_next_segment the
|
||||||
|
# count only drops to 1, so the fetch runs with the lock held.
|
||||||
self._buffer_lock = threading.RLock()
|
self._buffer_lock = threading.RLock()
|
||||||
|
|
||||||
# Plugin rotation, and the position of the next plugin to fetch in it.
|
# Plugin rotation, and the position of the next plugin to fetch in it.
|
||||||
@@ -536,7 +539,10 @@ class StreamManager:
|
|||||||
|
|
||||||
plugin_id = self._ordered_plugins[self._prefetch_index]
|
plugin_id = self._ordered_plugins[self._prefetch_index]
|
||||||
|
|
||||||
# Release lock for potentially slow content fetch
|
# Release for the potentially slow content fetch. This frees
|
||||||
|
# the lock only when the caller did not hold it already
|
||||||
|
# (initialize); under get_next_segment's hold the RLock count
|
||||||
|
# just drops to 1 and other threads still wait.
|
||||||
self._buffer_lock.release()
|
self._buffer_lock.release()
|
||||||
try:
|
try:
|
||||||
segment = self._fetch_plugin_content(plugin_id)
|
segment = self._fetch_plugin_content(plugin_id)
|
||||||
@@ -765,7 +771,6 @@ class StreamManager:
|
|||||||
continue
|
continue
|
||||||
if images:
|
if images:
|
||||||
self.stats['segments_fetched'] += 1
|
self.stats['segments_fetched'] += 1
|
||||||
if images:
|
|
||||||
group.append((plugin_id, images))
|
group.append((plugin_id, images))
|
||||||
else:
|
else:
|
||||||
group.append((plugin_id, None if defer_empty else []))
|
group.append((plugin_id, None if defer_empty else []))
|
||||||
|
|||||||
@@ -0,0 +1,93 @@
|
|||||||
|
"""BDF faces must not be shared between threads.
|
||||||
|
|
||||||
|
src/common/bdf_font.load_bdf_face caches ``freetype.Face`` objects per thread
|
||||||
|
because FreeType does not allow two threads to use one face at once
|
||||||
|
(``load_char`` rewrites its glyph slot). FontManager.font_cache and
|
||||||
|
element_style._font_cache are process-wide, and used to cache the returned
|
||||||
|
face for every thread -- so the display thread and a plugin update thread
|
||||||
|
drew through the same Face.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import threading
|
||||||
|
|
||||||
|
import freetype
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from src import element_style
|
||||||
|
from src.font_manager import FontManager
|
||||||
|
|
||||||
|
|
||||||
|
def _on_other_thread(fn):
|
||||||
|
box = {}
|
||||||
|
|
||||||
|
def run():
|
||||||
|
box['value'] = fn()
|
||||||
|
|
||||||
|
t = threading.Thread(target=run)
|
||||||
|
t.start()
|
||||||
|
t.join(timeout=10)
|
||||||
|
assert not t.is_alive()
|
||||||
|
return box['value']
|
||||||
|
|
||||||
|
|
||||||
|
def test_font_manager_gives_each_thread_its_own_bdf_face():
|
||||||
|
fm = FontManager({})
|
||||||
|
here = fm.get_font("five_by_seven", 7)
|
||||||
|
there = _on_other_thread(lambda: fm.get_font("five_by_seven", 7))
|
||||||
|
assert isinstance(here, freetype.Face) and isinstance(there, freetype.Face)
|
||||||
|
assert here is not there
|
||||||
|
# Within one thread the face is still reused.
|
||||||
|
assert fm.get_font("five_by_seven", 7) is here
|
||||||
|
|
||||||
|
|
||||||
|
def test_font_manager_still_caches_ttf():
|
||||||
|
fm = FontManager({})
|
||||||
|
here = fm.get_font("press_start", 8)
|
||||||
|
assert _on_other_thread(lambda: fm.get_font("press_start", 8)) is here
|
||||||
|
|
||||||
|
|
||||||
|
def test_element_style_gives_each_thread_its_own_bdf_face():
|
||||||
|
element_style._font_cache.clear()
|
||||||
|
try:
|
||||||
|
here = element_style.load_font("5x7.bdf", 7)
|
||||||
|
there = _on_other_thread(lambda: element_style.load_font("5x7.bdf", 7))
|
||||||
|
assert isinstance(here, freetype.Face) and isinstance(there, freetype.Face)
|
||||||
|
assert here is not there
|
||||||
|
assert element_style.load_font("5x7.bdf", 7) is here
|
||||||
|
finally:
|
||||||
|
element_style._font_cache.clear()
|
||||||
|
|
||||||
|
|
||||||
|
def test_element_style_cache_survives_concurrent_eviction(monkeypatch):
|
||||||
|
"""get() then move_to_end() on the shared OrderedDict raised KeyError when
|
||||||
|
another thread evicted the key in between. Forced deterministically: the
|
||||||
|
first get() hands the cache to a second thread that fills it past the
|
||||||
|
bound before get() returns."""
|
||||||
|
from collections import OrderedDict
|
||||||
|
|
||||||
|
owner = threading.get_ident()
|
||||||
|
state = {'worker': None}
|
||||||
|
|
||||||
|
class RacingCache(OrderedDict):
|
||||||
|
def get(self, key, default=None):
|
||||||
|
value = super().get(key, default)
|
||||||
|
if state['worker'] is None and threading.get_ident() == owner:
|
||||||
|
state['worker'] = threading.Thread(
|
||||||
|
target=lambda: [element_style.load_font(
|
||||||
|
"PressStart2P-Regular.ttf", s) for s in (20, 21, 22)],
|
||||||
|
daemon=True)
|
||||||
|
state['worker'].start()
|
||||||
|
# Unguarded, the worker evicts `key` now; guarded, it waits
|
||||||
|
# for the lock this thread holds.
|
||||||
|
state['worker'].join(timeout=0.5)
|
||||||
|
return value
|
||||||
|
|
||||||
|
monkeypatch.setattr(element_style, "_FONT_CACHE_MAX", 2)
|
||||||
|
cache = RacingCache()
|
||||||
|
monkeypatch.setattr(element_style, "_font_cache", cache)
|
||||||
|
element_style.load_font("PressStart2P-Regular.ttf", 8) # populate
|
||||||
|
state['worker'] = None
|
||||||
|
element_style.load_font("PressStart2P-Regular.ttf", 8) # hit, raced
|
||||||
|
state['worker'].join(timeout=10)
|
||||||
|
assert not state['worker'].is_alive()
|
||||||
|
assert len(cache) <= 2
|
||||||
@@ -0,0 +1,111 @@
|
|||||||
|
"""DisplayManager's deferred-update queue across threads.
|
||||||
|
|
||||||
|
defer_update() is called from plugin update() on the update worker thread
|
||||||
|
while process_deferred_updates() runs on the render thread. Both rebuild the
|
||||||
|
queue list (TTL filter, [n:] slice) and assign it back, so an append that
|
||||||
|
landed between one side's read and its assignment used to be dropped.
|
||||||
|
|
||||||
|
The race is forced deterministically: _scrolling_state is swapped for a dict
|
||||||
|
whose first store of 'deferred_updates' (made by the render side) first lets
|
||||||
|
a worker-thread defer_update() run. Unlocked, the worker appends to the list
|
||||||
|
the render side is about to overwrite; locked, the worker waits and appends
|
||||||
|
to the list that was stored.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import os
|
||||||
|
import sys
|
||||||
|
import threading
|
||||||
|
|
||||||
|
os.environ["EMULATOR"] = "true"
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
sys.path.insert(0, os.path.join(os.path.dirname(__file__), ".."))
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture(scope="module")
|
||||||
|
def dm():
|
||||||
|
from src.display_manager import DisplayManager
|
||||||
|
DisplayManager._instance = None
|
||||||
|
manager = DisplayManager({
|
||||||
|
"display": {
|
||||||
|
"hardware": {"rows": 32, "cols": 64, "chain_length": 1,
|
||||||
|
"parallel": 1, "brightness": 90},
|
||||||
|
"runtime": {"gpio_slowdown": 0},
|
||||||
|
},
|
||||||
|
}, suppress_test_pattern=True)
|
||||||
|
yield manager
|
||||||
|
DisplayManager._instance = None
|
||||||
|
|
||||||
|
|
||||||
|
class _InterleavingState(dict):
|
||||||
|
"""On the first 'deferred_updates' store from the owning thread, run
|
||||||
|
``interleave`` on another thread before the store happens."""
|
||||||
|
|
||||||
|
def __init__(self, base, interleave):
|
||||||
|
super().__init__(base)
|
||||||
|
dict.__setitem__(self, 'deferred_updates', [])
|
||||||
|
self._owner = threading.get_ident()
|
||||||
|
self._interleave = interleave
|
||||||
|
self.worker = None
|
||||||
|
|
||||||
|
def __setitem__(self, key, value):
|
||||||
|
if (key == 'deferred_updates' and self.worker is None
|
||||||
|
and threading.get_ident() == self._owner):
|
||||||
|
self.worker = threading.Thread(target=self._interleave, daemon=True)
|
||||||
|
self.worker.start()
|
||||||
|
# Unlocked code lets the worker finish here; locked code blocks
|
||||||
|
# it on the lock this thread holds, so give up waiting quickly.
|
||||||
|
self.worker.join(timeout=0.5)
|
||||||
|
super().__setitem__(key, value)
|
||||||
|
|
||||||
|
|
||||||
|
def test_defer_from_another_thread_is_not_lost(dm):
|
||||||
|
original = dm._scrolling_state
|
||||||
|
ran = [] # the callables never run here; they only need to exist
|
||||||
|
try:
|
||||||
|
state = _InterleavingState(
|
||||||
|
original, lambda: dm.defer_update(lambda: ran.append('late')))
|
||||||
|
dm._scrolling_state = state
|
||||||
|
dm.defer_update(lambda: ran.append('first'))
|
||||||
|
state.worker.join(timeout=5)
|
||||||
|
assert not state.worker.is_alive()
|
||||||
|
assert len(state['deferred_updates']) == 2, (
|
||||||
|
"a defer_update() from another thread was lost")
|
||||||
|
finally:
|
||||||
|
dm._scrolling_state = original
|
||||||
|
original['deferred_updates'] = []
|
||||||
|
|
||||||
|
|
||||||
|
def test_process_during_defer_keeps_both(dm):
|
||||||
|
original = dm._scrolling_state
|
||||||
|
try:
|
||||||
|
state = _InterleavingState(
|
||||||
|
original, lambda: dm.defer_update(lambda: None))
|
||||||
|
dict.__setitem__(state, 'is_scrolling', True)
|
||||||
|
dict.__setitem__(state, 'last_scroll_activity', 1e18) # stays "scrolling"
|
||||||
|
dm._scrolling_state = state
|
||||||
|
# Render side: only the TTL cleanup runs while scrolling.
|
||||||
|
dm.process_deferred_updates()
|
||||||
|
state.worker.join(timeout=5)
|
||||||
|
assert not state.worker.is_alive()
|
||||||
|
assert len(state['deferred_updates']) == 1, (
|
||||||
|
"a defer_update() during the render-side cleanup was lost")
|
||||||
|
finally:
|
||||||
|
dm._scrolling_state = original
|
||||||
|
original['deferred_updates'] = []
|
||||||
|
|
||||||
|
|
||||||
|
def test_queued_callable_may_defer_again(dm):
|
||||||
|
"""The callables run outside the lock, so one that re-defers (a plugin
|
||||||
|
retrying later) must not deadlock the render thread."""
|
||||||
|
dm._scrolling_state['is_scrolling'] = False
|
||||||
|
try:
|
||||||
|
dm.defer_update(lambda: dm.defer_update(lambda: None))
|
||||||
|
t = threading.Thread(target=dm.process_deferred_updates, daemon=True)
|
||||||
|
t.start()
|
||||||
|
t.join(timeout=5)
|
||||||
|
assert not t.is_alive(), "process_deferred_updates deadlocked"
|
||||||
|
assert len(dm._scrolling_state['deferred_updates']) == 1
|
||||||
|
finally:
|
||||||
|
dm._scrolling_state['deferred_updates'] = []
|
||||||
@@ -0,0 +1,31 @@
|
|||||||
|
"""Errors DisplayController deliberately swallows must still leave a trace.
|
||||||
|
|
||||||
|
Several ``except Exception: pass`` blocks hid plugin and filesystem faults
|
||||||
|
completely; they now log at DEBUG with the traceback, and still don't raise.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import logging
|
||||||
|
import os
|
||||||
|
from unittest.mock import MagicMock
|
||||||
|
|
||||||
|
os.environ.setdefault("EMULATOR", "true")
|
||||||
|
|
||||||
|
from src.display_controller import DisplayController # noqa: E402
|
||||||
|
|
||||||
|
|
||||||
|
def test_has_live_content_failure_is_logged_not_raised(caplog):
|
||||||
|
dc = object.__new__(DisplayController)
|
||||||
|
broken = MagicMock()
|
||||||
|
broken.has_live_content.side_effect = RuntimeError("feed parse failed")
|
||||||
|
dc.plugin_display_modes = {"nfl": ["nfl_live", "nfl_recent"]}
|
||||||
|
dc.mode_to_plugin_id = {}
|
||||||
|
dc.plugin_modes = {"nfl_live": broken, "nfl_recent": MagicMock()}
|
||||||
|
|
||||||
|
with caplog.at_level(logging.DEBUG, logger="src.display_controller"):
|
||||||
|
modes = dc._on_demand_modes_for_plugin("nfl")
|
||||||
|
|
||||||
|
# Behaviour unchanged: the raising live mode counts as having no content.
|
||||||
|
assert modes == ["nfl_recent"]
|
||||||
|
records = [r for r in caplog.records
|
||||||
|
if "has_live_content() failed for nfl_live" in r.getMessage()]
|
||||||
|
assert len(records) == 1 and records[0].exc_info is not None
|
||||||
@@ -47,3 +47,28 @@ def test_debug_output_appears_when_the_root_is_at_debug():
|
|||||||
assert dm.logger.isEnabledFor(logging.DEBUG)
|
assert dm.logger.isEnabledFor(logging.DEBUG)
|
||||||
finally:
|
finally:
|
||||||
root.setLevel(previous)
|
root.setLevel(previous)
|
||||||
|
|
||||||
|
|
||||||
|
def test_update_display_errors_are_rate_limited(caplog):
|
||||||
|
# update_display() runs every frame; a persistent fault logged an ERROR
|
||||||
|
# line per frame (~100 a second) and never a traceback.
|
||||||
|
from unittest.mock import MagicMock
|
||||||
|
|
||||||
|
dm_obj = object.__new__(dm.DisplayManager)
|
||||||
|
dm_obj._writes_suppressed = MagicMock(side_effect=RuntimeError("boom"))
|
||||||
|
|
||||||
|
with caplog.at_level(logging.ERROR, logger='src.display_manager'):
|
||||||
|
for _ in range(50):
|
||||||
|
dm_obj.update_display() # must not raise
|
||||||
|
errors = [r for r in caplog.records
|
||||||
|
if r.getMessage().startswith('Error updating display')]
|
||||||
|
assert len(errors) == 1
|
||||||
|
assert errors[0].exc_info is not None
|
||||||
|
|
||||||
|
# Once the interval has passed, one more line reports what was skipped.
|
||||||
|
dm_obj._update_error_logged_at -= dm._UPDATE_ERROR_LOG_INTERVAL + 1
|
||||||
|
dm_obj.update_display()
|
||||||
|
errors = [r for r in caplog.records
|
||||||
|
if r.getMessage().startswith('Error updating display')]
|
||||||
|
assert len(errors) == 2
|
||||||
|
assert '49 more' in errors[1].getMessage()
|
||||||
|
|||||||
@@ -0,0 +1,78 @@
|
|||||||
|
"""Follower adoption of the leader's scroll image (src/display_controller.py).
|
||||||
|
|
||||||
|
The leader's image arrives on the sync TCP thread. That callback used to set
|
||||||
|
scroll_helper.cached_image, cached_array and total_scroll_width one after
|
||||||
|
another while the render thread sliced frames out of them, so a frame could
|
||||||
|
pair the new array with the old width. The callback now only queues the
|
||||||
|
image; the render thread swaps all three in at the start of a follower frame.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import os
|
||||||
|
from collections import deque
|
||||||
|
from types import SimpleNamespace
|
||||||
|
from unittest.mock import MagicMock, patch
|
||||||
|
|
||||||
|
os.environ.setdefault("EMULATOR", "true")
|
||||||
|
|
||||||
|
from PIL import Image # noqa: E402
|
||||||
|
|
||||||
|
from src.display_controller import DisplayController # noqa: E402
|
||||||
|
|
||||||
|
|
||||||
|
def _wired_controller():
|
||||||
|
dc = object.__new__(DisplayController)
|
||||||
|
dc.config = {"display": {"vegas_scroll": {"enabled": True}}, "sync": {}}
|
||||||
|
dc.display_manager = MagicMock()
|
||||||
|
dc.plugin_manager = MagicMock()
|
||||||
|
dc.sync_manager = MagicMock()
|
||||||
|
dc._check_live_priority = MagicMock()
|
||||||
|
dc._check_vegas_interrupt = MagicMock(return_value=False)
|
||||||
|
dc._follower_pending_new_image = True
|
||||||
|
dc._follower_incoming_image = deque(maxlen=1)
|
||||||
|
|
||||||
|
old = Image.new("RGB", (100, 16))
|
||||||
|
helper = SimpleNamespace(cached_image=old, cached_array="old-array",
|
||||||
|
total_scroll_width=100)
|
||||||
|
coordinator = MagicMock()
|
||||||
|
coordinator.render_pipeline = SimpleNamespace(scroll_helper=helper)
|
||||||
|
with patch('src.vegas_mode.VegasModeCoordinator',
|
||||||
|
MagicMock(return_value=coordinator)):
|
||||||
|
dc._initialize_vegas_mode()
|
||||||
|
on_image = dc.sync_manager.set_on_scroll_image.call_args[0][0]
|
||||||
|
return dc, coordinator.render_pipeline, helper, old, on_image
|
||||||
|
|
||||||
|
|
||||||
|
def test_tcp_callback_does_not_touch_the_render_state():
|
||||||
|
dc, _rp, helper, old, on_image = _wired_controller()
|
||||||
|
on_image(Image.new("RGB", (300, 16)))
|
||||||
|
# Nothing the render thread reads has changed yet.
|
||||||
|
assert helper.cached_image is old
|
||||||
|
assert helper.cached_array == "old-array"
|
||||||
|
assert helper.total_scroll_width == 100
|
||||||
|
assert dc._follower_pending_new_image is True
|
||||||
|
|
||||||
|
|
||||||
|
def test_render_thread_adopts_all_three_together():
|
||||||
|
dc, rp, helper, _old, on_image = _wired_controller()
|
||||||
|
new = Image.new("RGB", (300, 16), (255, 0, 0))
|
||||||
|
on_image(new)
|
||||||
|
dc._adopt_follower_scroll_image(rp)
|
||||||
|
assert helper.cached_image is new
|
||||||
|
assert helper.cached_array.shape == (16, 300, 3)
|
||||||
|
assert helper.cached_array[0, 0].tolist() == [255, 0, 0]
|
||||||
|
assert helper.total_scroll_width == 300
|
||||||
|
assert dc._follower_pending_new_image is False
|
||||||
|
# Consumed: a second frame changes nothing.
|
||||||
|
helper.total_scroll_width = 7
|
||||||
|
dc._adopt_follower_scroll_image(rp)
|
||||||
|
assert helper.total_scroll_width == 7
|
||||||
|
|
||||||
|
|
||||||
|
def test_only_the_latest_image_is_adopted():
|
||||||
|
dc, rp, helper, _old, on_image = _wired_controller()
|
||||||
|
on_image(Image.new("RGB", (200, 16)))
|
||||||
|
latest = Image.new("RGB", (400, 16))
|
||||||
|
on_image(latest)
|
||||||
|
dc._adopt_follower_scroll_image(rp)
|
||||||
|
assert helper.cached_image is latest
|
||||||
|
assert helper.total_scroll_width == 400
|
||||||
@@ -122,4 +122,24 @@ def test_display_manager_defaults_match_display_manager():
|
|||||||
for field, value in ms.DISPLAY_MANAGER_DEFAULTS.items():
|
for field, value in ms.DISPLAY_MANAGER_DEFAULTS.items():
|
||||||
if field in ('rows', 'cols', 'chain_length', 'parallel'):
|
if field in ('rows', 'cols', 'chain_length', 'parallel'):
|
||||||
continue # read through display_geometry's DEFAULT_* constants
|
continue # read through display_geometry's DEFAULT_* constants
|
||||||
|
if field == 'limit_refresh_rate_hz':
|
||||||
|
assert "get('limit_refresh_rate_hz', DEFAULT_REFRESH_LIMIT_HZ)" in source
|
||||||
|
continue
|
||||||
assert f"get('{field}', {value!r})" in source, field
|
assert f"get('{field}', {value!r})" in source, field
|
||||||
|
|
||||||
|
|
||||||
|
def test_missing_refresh_limit_applies_the_rate_pacing_assumes():
|
||||||
|
"""With limit_refresh_rate_hz absent the library was capped at 90 Hz while
|
||||||
|
refresh_hz (what scroll pacing solves against) reported 100. One default,
|
||||||
|
and it is the template's."""
|
||||||
|
import os
|
||||||
|
from types import SimpleNamespace
|
||||||
|
os.environ.setdefault("EMULATOR", "true") # import off-Pi, as other modules do
|
||||||
|
from src.display_manager import DisplayManager
|
||||||
|
|
||||||
|
options = DisplayManager.apply_matrix_options(SimpleNamespace(), {})
|
||||||
|
reported = DisplayManager.refresh_hz.fget(SimpleNamespace(config={}))
|
||||||
|
assert options.limit_refresh_rate_hz == reported
|
||||||
|
template = json.loads((REPO_ROOT / 'config' / 'config.template.json').read_text(encoding='utf-8'))
|
||||||
|
assert options.limit_refresh_rate_hz == template['display']['hardware']['limit_refresh_rate_hz']
|
||||||
|
assert ms.DISPLAY_MANAGER_DEFAULTS['limit_refresh_rate_hz'] == options.limit_refresh_rate_hz
|
||||||
|
|||||||
@@ -0,0 +1,21 @@
|
|||||||
|
"""The repo-root display_controller.py goes through run.py.
|
||||||
|
|
||||||
|
It used to call src.display_controller.main() directly, skipping run.py's
|
||||||
|
sys.dont_write_bytecode, its -e/-d flags and its logging setup. run.py's
|
||||||
|
argument parser answering --help shows the shim now runs run.py.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import subprocess
|
||||||
|
import sys
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
REPO_ROOT = Path(__file__).resolve().parents[1]
|
||||||
|
|
||||||
|
|
||||||
|
def test_root_display_controller_runs_run_py():
|
||||||
|
result = subprocess.run(
|
||||||
|
[sys.executable, str(REPO_ROOT / "display_controller.py"), "--help"],
|
||||||
|
cwd=str(REPO_ROOT), capture_output=True, text=True, timeout=60,
|
||||||
|
)
|
||||||
|
assert result.returncode == 0, result.stderr
|
||||||
|
assert "--emulator" in result.stdout and "--debug" in result.stdout
|
||||||
Reference in New Issue
Block a user