Merge branch 'claude/hdpi-scroll-performance-antialiasing-4ae609' into claude/offscreen-rendering

# Conflicts:
#	src/display_manager.py
#	src/vegas_mode/config.py
#	src/vegas_mode/render_pipeline.py
This commit is contained in:
Chuck
2026-09-24 17:46:58 -04:00
89 changed files with 2887 additions and 5491 deletions
+2 -2
View File
@@ -6,8 +6,8 @@ plugins' content is composed into a single horizontally scrolling display.
Components:
- VegasModeCoordinator: Main orchestrator for Vegas mode
- StreamManager: Manages plugin content streaming with 1-2 ahead buffering
- RenderPipeline: Handles 125 FPS rendering with double-buffering
- StreamManager: Plugin rotation, content fetching and pending-update tracking
- RenderPipeline: Strip composition and per-frame rendering
- PluginAdapter: Converts plugin content to scrollable images
- VegasModeConfig: Configuration management
"""
+42 -52
View File
@@ -201,7 +201,7 @@ class VegasModeConfig:
# Dynamic duration
dynamic_duration_enabled: bool = True
min_cycle_duration: int = 60 # Minimum seconds per full cycle
max_cycle_duration: int = 600 # Maximum seconds per full cycle
max_cycle_duration: int = 240 # Maximum seconds per full cycle
@classmethod
def from_config(cls, config: Dict[str, Any]) -> 'VegasModeConfig':
@@ -215,48 +215,54 @@ class VegasModeConfig:
VegasModeConfig instance
"""
vegas_config = config.get('display', {}).get('vegas_scroll', {})
# Missing keys fall back to the field defaults above, so each default
# is written once and the two cannot drift apart.
d = cls()
get = vegas_config.get
return cls(
enabled=vegas_config.get('enabled', False),
scroll_speed=float(vegas_config.get('scroll_speed', 50.0)),
separator_width=int(vegas_config.get('separator_width', 32)),
intra_plugin_gap=int(vegas_config.get('intra_plugin_gap', 8)),
render_width_pct=int(vegas_config.get('render_width_pct', 100)),
enabled=get('enabled', d.enabled),
scroll_speed=float(get('scroll_speed', d.scroll_speed)),
separator_width=int(get('separator_width', d.separator_width)),
intra_plugin_gap=int(get('intra_plugin_gap', d.intra_plugin_gap)),
render_width_pct=int(get('render_width_pct', d.render_width_pct)),
min_content_separation=int(
vegas_config.get('min_content_separation', 24)),
min_cut_gap=int(vegas_config.get('min_cut_gap', 6)),
smooth_scroll=vegas_config.get('smooth_scroll', True),
sub_pixel_blend=bool(vegas_config.get('sub_pixel_blend', False)),
continuous_scroll=vegas_config.get('continuous_scroll', True),
offscreen_prefetch=bool(vegas_config.get('offscreen_prefetch', True)),
switch_interval_ms=float(vegas_config.get('switch_interval_ms', 0.0) or 0.0),
prefetch_gate=bool(vegas_config.get('prefetch_gate', True)),
get('min_content_separation', d.min_content_separation)),
min_cut_gap=int(get('min_cut_gap', d.min_cut_gap)),
smooth_scroll=get('smooth_scroll', d.smooth_scroll),
sub_pixel_blend=bool(get('sub_pixel_blend', d.sub_pixel_blend)),
continuous_scroll=get('continuous_scroll', d.continuous_scroll),
offscreen_prefetch=bool(get('offscreen_prefetch', d.offscreen_prefetch)),
switch_interval_ms=float(get('switch_interval_ms', d.switch_interval_ms) or 0.0),
prefetch_gate=bool(get('prefetch_gate', d.prefetch_gate)),
extend_threshold_screens=float(
vegas_config.get('extend_threshold_screens', 2.0)),
auto_trim=vegas_config.get('auto_trim', True),
trim_threshold=int(vegas_config.get('trim_threshold', 10)),
content_padding=int(vegas_config.get('content_padding', 8)),
min_plugin_width=int(vegas_config.get('min_plugin_width', 8)),
lead_in_width=int(vegas_config.get('lead_in_width', 0)),
plugins_per_cycle=int(vegas_config.get('plugins_per_cycle', 6)),
get('extend_threshold_screens', d.extend_threshold_screens)),
auto_trim=get('auto_trim', d.auto_trim),
trim_threshold=int(get('trim_threshold', d.trim_threshold)),
content_padding=int(get('content_padding', d.content_padding)),
min_plugin_width=int(get('min_plugin_width', d.min_plugin_width)),
lead_in_width=int(get('lead_in_width', d.lead_in_width)),
plugins_per_cycle=int(get('plugins_per_cycle', d.plugins_per_cycle)),
max_plugin_width_ratio=float(
vegas_config.get('max_plugin_width_ratio', 0.0)),
overflow_mode=str(vegas_config.get('overflow_mode', 'rotate')),
plugin_order=list(vegas_config.get('plugin_order', [])),
excluded_plugins=set(vegas_config.get('excluded_plugins', [])),
live_in_ticker=bool(vegas_config.get('live_in_ticker', False)),
get('max_plugin_width_ratio', d.max_plugin_width_ratio)),
overflow_mode=str(get('overflow_mode', d.overflow_mode)),
plugin_order=list(get('plugin_order', d.plugin_order)),
excluded_plugins=set(get('excluded_plugins', d.excluded_plugins)),
live_in_ticker=bool(get('live_in_ticker', d.live_in_ticker)),
# Clamped: a weight below 1 would drop the plugin from the rotation
# entirely, and a very large one starves everything else.
live_weight=max(1, min(10, int(vegas_config.get('live_weight', 3)))),
favorite_live_weight=max(
1, min(10, int(vegas_config.get('favorite_live_weight', 5)))),
target_fps=int(vegas_config.get('target_fps', 125)),
buffer_ahead=int(vegas_config.get('buffer_ahead', 2)),
frame_based_scrolling=vegas_config.get('frame_based_scrolling', True),
scroll_delay=float(vegas_config.get('scroll_delay', 0.02)),
dynamic_duration_enabled=vegas_config.get('dynamic_duration_enabled', True),
min_cycle_duration=int(vegas_config.get('min_cycle_duration', 60)),
max_cycle_duration=int(vegas_config.get('max_cycle_duration', 600)),
live_weight=max(1, min(10, int(get('live_weight', d.live_weight)))),
favorite_live_weight=max(1, min(10, int(
get('favorite_live_weight', d.favorite_live_weight)))),
target_fps=int(get('target_fps', d.target_fps)),
buffer_ahead=int(get('buffer_ahead', d.buffer_ahead)),
frame_based_scrolling=get(
'frame_based_scrolling', d.frame_based_scrolling),
scroll_delay=float(get('scroll_delay', d.scroll_delay)),
dynamic_duration_enabled=get(
'dynamic_duration_enabled', d.dynamic_duration_enabled),
min_cycle_duration=int(get('min_cycle_duration', d.min_cycle_duration)),
max_cycle_duration=int(get('max_cycle_duration', d.max_cycle_duration)),
)
def to_dict(self) -> Dict[str, Any]:
@@ -302,22 +308,6 @@ class VegasModeConfig:
"""Get the frame interval in seconds for target FPS."""
return 1.0 / max(1, self.target_fps)
def is_plugin_included(self, plugin_id: str) -> bool:
"""
Check if a plugin should be included in Vegas scroll.
This is consistent with get_ordered_plugins - plugins not explicitly
in plugin_order are still included (appended at the end) unless excluded.
Args:
plugin_id: Plugin identifier to check
Returns:
True if plugin should be included
"""
# Plugins are included unless explicitly excluded
return plugin_id not in self.excluded_plugins
def get_ordered_plugins(self, available_plugins: List[str]) -> List[str]:
"""
Get plugins in configured order, filtering excluded ones.
+5 -45
View File
@@ -148,13 +148,8 @@ class VegasModeCoordinator:
# Static pause handling
self._static_pause_active = False
self._static_pause_plugin: Optional['BasePlugin'] = None
self._static_pause_start: Optional[float] = None
self._saved_scroll_position: Optional[int] = None
# Track which plugins should use STATIC mode (pause scroll)
self._static_mode_plugins: set = set()
# Statistics
self.stats = {
'total_runtime_seconds': 0.0,
@@ -238,7 +233,9 @@ class VegasModeCoordinator:
returns immediately, collapsing the inter-iteration gap to <1 ms.
Args:
callback: Callable with no arguments (typically _tick_plugin_updates)
callback: Callable with no arguments. The display controller
passes _tick_plugin_updates_for_vegas, which also reports the
plugins that got fresh data through mark_plugin_updated().
"""
self._update_callback = callback
@@ -476,9 +473,6 @@ class VegasModeCoordinator:
if not self.start():
return False
# Update static mode plugin list on iteration start
self._update_static_mode_plugins()
if self.vegas_config.continuous_scroll:
# The strip is continuously extended and trimmed, so its width says
# nothing about how long to run. This is only how often control
@@ -798,13 +792,6 @@ class VegasModeCoordinator:
return status
def get_ordered_plugins(self) -> List[str]:
"""Get the current ordered list of plugins in Vegas scroll."""
if hasattr(self.plugin_manager, 'plugins'):
available = list(self.plugin_manager.plugins.keys())
return self.vegas_config.get_ordered_plugins(available)
return []
# -------------------------------------------------------------------------
# Static pause handling (for STATIC display mode)
# -------------------------------------------------------------------------
@@ -860,8 +847,6 @@ class VegasModeCoordinator:
# Save current scroll position for smooth resume
self._saved_scroll_position = self.render_pipeline.get_scroll_position()
self._static_pause_active = True
self._static_pause_plugin = plugin
self._static_pause_start = time.time()
self.stats['static_pauses'] += 1
logger.info("Static pause started for plugin: %s", plugin_id)
@@ -888,9 +873,9 @@ class VegasModeCoordinator:
logger.info("Static pause interrupted by live priority")
return False
# Yield immediately if multi-display follower mode becomes active
# On-demand, a WiFi message, the schedule, follower mode...
if self._interrupt_check and self._interrupt_check():
logger.info("Static pause interrupted by sync follower mode")
logger.info("Static pause interrupted by the display controller")
return False
# Sleep in small increments to remain responsive
@@ -925,8 +910,6 @@ class VegasModeCoordinator:
# Clear pause state
self._static_pause_active = False
self._static_pause_plugin = None
self._static_pause_start = None
# Restore scroll position if we're resuming
if should_resume_scrolling and self._saved_scroll_position is not None:
@@ -940,29 +923,6 @@ class VegasModeCoordinator:
else:
logger.debug("Static pause ended (interrupted, not resuming scroll)")
def _update_static_mode_plugins(self) -> None:
"""Update the set of plugins using STATIC display mode."""
self._static_mode_plugins.clear()
for plugin_id in self.get_ordered_plugins():
plugin = self.plugin_manager.get_plugin(plugin_id)
if plugin:
try:
mode = plugin.get_vegas_display_mode()
if mode == VegasDisplayMode.STATIC:
self._static_mode_plugins.add(plugin_id)
except Exception:
logger.exception(
"Error getting vegas display mode for plugin %s",
plugin_id
)
if self._static_mode_plugins:
logger.info(
"Static mode plugins: %s",
', '.join(self._static_mode_plugins)
)
def cleanup(self) -> None:
"""Clean up all resources."""
self.stop()
-48
View File
@@ -230,54 +230,6 @@ def blank_runs(
return list(zip(starts[long_enough].tolist(), ends[long_enough].tolist()))
def find_blank_cut(
img: Image.Image,
target: int,
search_radius: int,
threshold: int = DEFAULT_INK_THRESHOLD,
) -> int:
"""
Find a column near ``target`` that carries no ink, so an image can be cut
there without slicing through a glyph or logo.
Used when a single oversized segment has to be narrowed to fit a width
budget. Cutting at an arbitrary column would leave half a character
hanging at the panel edge; snapping to the nearest gap hides the cut.
Args:
img: Image to cut
target: Preferred cut column
search_radius: How far either side of ``target`` to look
threshold: Ink threshold
Returns:
A blank column within the search window, or ``target`` clamped to the
image bounds when the window contains no blank column at all.
"""
width = img.width
target = max(0, min(target, width))
if search_radius <= 0 or width == 0:
return target
ink = column_has_ink(img, threshold)
# target may legitimately equal width (a cut after the last column), but
# there is no column to inspect there, so both bounds stop at width - 1.
lo = max(0, min(target - search_radius, width - 1))
hi = max(0, min(target + search_radius, width - 1))
# Walk outwards from target so the nearest gap wins.
for offset in range(0, search_radius + 1):
right = target + offset
if lo <= right <= hi and not ink[right]:
return right
left = target - offset
if lo <= left <= hi and not ink[left]:
return left
return target
class DeadWindowStats(NamedTuple):
"""How much of a composed ticker reads as blank to a viewer."""
+14 -17
View File
@@ -59,15 +59,8 @@ class PluginAdapter:
from src.vegas_mode.config import VegasModeConfig
config = VegasModeConfig()
self.config = config
# Handle both property and method access patterns
self.display_width = (
display_manager.width() if callable(display_manager.width)
else display_manager.width
)
self.display_height = (
display_manager.height() if callable(display_manager.height)
else display_manager.height
)
self.display_width = display_manager.width
self.display_height = display_manager.height
# Cache for recently fetched content (prevents redundant fetch)
self._content_cache: dict = {}
@@ -263,13 +256,14 @@ class PluginAdapter:
Trim dead space off a segment, then cache it.
Every content path funnels through here so trimming is applied
uniformly. Previously only the scroll_helper path had its margins
stripped, which left plugins that render onto a full-display canvas
contributing their entire blank canvas to the ticker.
uniformly; a plugin that renders onto a full-display canvas would
otherwise contribute its whole blank canvas to the ticker.
Each image is trimmed independently because compose_scroll_content()
treats every image as its own item and inserts separator_width between
them — so a per-image trim is what makes that separator the real gap.
Each image is trimmed independently. The render pipeline joins one
plugin's images with a gap measured from their ink
(RenderPipeline._join_plugin_rows) and puts separator_width only
between plugins, so the margins a row keeps are content_padding, not
whatever blank canvas the plugin happened to draw it on.
Args:
images: Raw content from one of the fetch paths
@@ -680,8 +674,11 @@ class PluginAdapter:
Narrow a single oversized image to the budget, advancing a window
through it across cycles.
The cut is snapped to the nearest blank column so it does not slice
through a glyph or logo and leave half a character at the panel edge.
Cuts land only at item boundaries: the middle of a blank run at least
``min_cut_gap`` columns wide. The window ends at the last boundary
inside the budget, or overruns to the next one when there is none, so
an item is never sliced. An image with no such runs (a map, a chart)
is continuous content and is cropped to the budget exactly.
Rotation is tracked as an index into the strip's item boundaries rather
than as a pixel column, because a ticker re-renders between fetches. A
+26 -61
View File
@@ -1,8 +1,8 @@
"""
Render Pipeline for Vegas Mode
Handles high-FPS (125 FPS) rendering with double-buffering for smooth scrolling.
Uses the existing ScrollHelper for numpy-optimized scroll operations.
Composes plugin content into one wide strip and renders the visible window of
it each frame, using ScrollHelper for the numpy-backed scroll.
"""
import logging
@@ -11,7 +11,7 @@ import time
import threading
from collections import deque
from contextlib import nullcontext
from typing import Optional, List, Any, Dict, Deque, TYPE_CHECKING
from typing import Optional, List, Any, Dict, Deque
from PIL import Image
from src.common.scroll_config import solve_crisp
@@ -20,11 +20,13 @@ from src.vegas_mode.config import VegasModeConfig
from src.vegas_mode.geometry import separation_gap
from src.vegas_mode.stream_manager import StreamManager
if TYPE_CHECKING:
pass
logger = logging.getLogger(__name__)
#: Shortest gap between multi-display sync sends from the leader, for both the
#: Vegas scroll position and the controller's per-frame follower images. The
#: payloads are raw and cheap, and 90/s is above the follower's render rate.
SYNC_SEND_INTERVAL = 1.0 / 90
class RenderPipeline:
"""
@@ -33,8 +35,9 @@ class RenderPipeline:
Key responsibilities:
- Compose content segments into scrollable image
- Manage scroll position and velocity
- Handle 125 FPS rendering loop
- Double-buffer for hot-swap during updates
- Render one frame per call at the target FPS
- Extend the strip (continuous mode) or recompose and hot-swap it (swap
mode) as content changes
- Track scroll cycle completion
"""
@@ -68,18 +71,10 @@ class RenderPipeline:
self.stream_manager = stream_manager
self.sync_manager = None # Optional DisplaySyncManager — set by coordinator
self.sync_follower_left = True # True = follower is LEFT of leader (default)
self._sync_send_interval = 1.0 / 90 # raw bytes are cheap; 90fps > follower render rate
self._last_sync_send = 0.0
# Display dimensions (handle both property and method access patterns)
self.display_width = (
display_manager.width() if callable(display_manager.width)
else display_manager.width
)
self.display_height = (
display_manager.height() if callable(display_manager.height)
else display_manager.height
)
self.display_width = display_manager.width
self.display_height = display_manager.height
# ScrollHelper for optimized scrolling
self.scroll_helper = ScrollHelper(
@@ -96,11 +91,6 @@ class RenderPipeline:
# Configure scroll helper
self._configure_scroll_helper()
# Double-buffer for composed images
self._active_scroll_image: Optional[Image.Image] = None
self._staging_scroll_image: Optional[Image.Image] = None
self._buffer_lock = threading.Lock()
# Group prepared off the render thread, waiting to be appended.
self._prepared_group = None
# Plugins that need the shared canvas, appended one at a time.
@@ -110,12 +100,10 @@ class RenderPipeline:
self._prefetch_lock = threading.Lock()
# Render state
self._is_rendering = False
self._cycle_complete = False
self._segments_in_scroll: List[str] = [] # Plugin IDs in current scroll
# Timing
self._last_frame_time = 0.0
# The sub-pixel path's pacing; the crisp path solves its own (frame_interval).
self._frame_interval = config.get_frame_interval()
self._cycle_start_time = 0.0
@@ -316,10 +304,6 @@ class RenderPipeline:
logger.error("ScrollHelper failed to create cached image")
return False
# Store reference to composed image
with self._buffer_lock:
self._active_scroll_image = self.scroll_helper.cached_image
# Track which plugins are in this scroll (get safely via buffer status)
self._segments_in_scroll = self.stream_manager.get_active_plugin_ids()
@@ -463,8 +447,6 @@ class RenderPipeline:
element_gap=0,
)
if appended:
with self._buffer_lock:
self._active_scroll_image = self.scroll_helper.cached_image
logger.info(
"[%s] Appended deferred content: strip now %dpx, %dpx ahead",
plugin_id, self.scroll_helper.total_scroll_width,
@@ -551,9 +533,6 @@ class RenderPipeline:
# Keep a screen's worth behind the viewport as a safety margin.
self.scroll_helper.drop_scrolled_prefix(keep_before=self.display_width)
with self._buffer_lock:
self._active_scroll_image = self.scroll_helper.cached_image
self._segments_in_scroll = [pid for pid, _ in grouped]
self.stats['composition_count'] += 1
self.stats['extensions'] = self.stats.get('extensions', 0) + 1
@@ -696,7 +675,7 @@ class RenderPipeline:
# leader's via TCP image transfer at each new_cycle) at scroll_x ± display_width.
if self.sync_manager:
now = time.time()
if now - self._last_sync_send >= self._sync_send_interval:
if now - self._last_sync_send >= SYNC_SEND_INTERVAL:
self._last_sync_send = now
self.sync_manager.send_scroll_x(self.scroll_helper.scroll_position)
@@ -735,7 +714,6 @@ class RenderPipeline:
Returns True when:
- Cycle is complete and we should start fresh
- Staging buffer has new content
- A plugin currently visible in the scroll has pending updated data
(e.g. a live score changed) — standalone (non-sync) mode only
"""
@@ -745,16 +723,11 @@ class RenderPipeline:
# When multi-display sync is active, defer mid-cycle hot swaps until the
# cycle ends naturally. Hot swaps block the render loop for 15-30ms while
# the image is rebuilt, causing a freeze+jump that the follower perceives
# as a speed-up. Deferring to cycle boundaries keeps transitions clean.
# Staging buffer content is still pre-loaded; it just applies at cycle end.
# as a speed-up. Deferring to cycle boundaries keeps transitions clean;
# the pending updates are still applied, by the recompose at cycle end.
if self.sync_manager is not None:
return False
# Check if we need more content in the buffer
buffer_status = self.stream_manager.get_buffer_status()
if buffer_status['staging_count'] > 0:
return True
# Trigger recompose when pending updates affect visible segments, so
# live score/status changes reach the display within a few seconds
# instead of waiting for the next full cycle.
@@ -784,10 +757,12 @@ class RenderPipeline:
def hot_swap_content(self) -> bool:
"""
Hot-swap to new composed content.
Refetch the plugins with pending updates and recompose the strip.
Called when staging buffer has updated content.
Swaps atomically to prevent visual glitches.
Swap mode only (continuous mode uses :meth:`refresh_updated_plugins`).
Called when :meth:`should_recompose` finds a visible plugin with
pending updates. The scroll resumes at the same relative position in
the rebuilt strip.
Returns:
True if swap occurred
@@ -800,9 +775,7 @@ class RenderPipeline:
old_width = self.scroll_helper.total_scroll_width
old_pos = self.scroll_helper.scroll_position
# Process any pending updates
self.stream_manager.process_updates()
self.stream_manager.swap_buffers()
# Recompose with updated content
if self.compose_scroll_content():
@@ -861,21 +834,17 @@ class RenderPipeline:
result = self.compose_scroll_content()
if result and self.sync_manager:
# When sync is active, start the leader past the lead-in gap so it
# immediately shows content, leaving the follower on the blank gap
# for a clean transition rather than near-end content wrapping
# around. This tracks lead_in_width rather than assuming a full
# display width of gap, which is no longer the default.
# Start the leader past the lead-in gap so it immediately shows
# content, leaving the follower on the blank gap for a clean
# transition rather than near-end content wrapping around.
self.scroll_helper.scroll_position = float(self.config.lead_in_width)
if result and self.sync_manager:
# Signal follower that a new cycle started (triggers its own rebuild)
self.sync_manager.send_new_cycle()
# Push the actual scroll image over TCP so follower has identical pixels.
# Done in a background thread to not block the render loop (~15ms transfer).
if self.scroll_helper.cached_image is not None:
import threading as _t
_t.Thread(
threading.Thread(
target=self.sync_manager.send_scroll_image,
args=(self.scroll_helper.cached_image,),
daemon=True, name="sync-image-push"
@@ -937,10 +906,6 @@ class RenderPipeline:
self.scroll_helper.reset_scroll()
self.scroll_helper.clear_cache()
with self._buffer_lock:
self._active_scroll_image = None
self._staging_scroll_image = None
self._cycle_complete = False
self._segments_in_scroll = []
self._frame_times = deque(maxlen=100)
+63 -140
View File
@@ -31,22 +31,14 @@ logger = logging.getLogger(__name__)
@dataclass
class ContentSegment:
"""Represents a segment of scrollable content from a plugin."""
"""One plugin's content for a cycle.
A STATIC segment carries no images: it marks where the coordinator pauses
the scroll to show the plugin full-screen.
"""
plugin_id: str
images: List[Image.Image]
total_width: int
display_mode: VegasDisplayMode = field(default=VegasDisplayMode.FIXED_SEGMENT)
fetched_at: float = field(default_factory=time.time)
is_stale: bool = False
@property
def image_count(self) -> int:
return len(self.images)
@property
def is_static(self) -> bool:
"""Check if this segment should trigger a static pause."""
return self.display_mode == VegasDisplayMode.STATIC
class StreamManager:
@@ -54,10 +46,12 @@ class StreamManager:
Manages streaming of plugin content for Vegas scroll mode.
Key responsibilities:
- Maintain ordered list of plugins to stream
- Prefetch content 1-2 plugins ahead of current position
- Handle plugin data updates via double-buffer swap
- Manage content lifecycle and staleness
- Maintain the ordered (and priority-weighted) rotation of plugins
- Fill the active buffer with a cycle's worth of segments (swap mode), or
hand out the next group of plugins directly (continuous mode)
- Track plugins whose data changed in ``_pending_updates``, which
:meth:`process_updates` (swap mode) or
:meth:`invalidate_pending_updates` (continuous mode) consumes
"""
def __init__(
@@ -78,14 +72,14 @@ class StreamManager:
self.plugin_manager = plugin_manager
self.plugin_adapter = plugin_adapter
# Content queue (double-buffered)
# Segments composed into the current cycle (swap mode only).
self._active_buffer: Deque[ContentSegment] = deque()
self._staging_buffer: Deque[ContentSegment] = deque()
self._buffer_lock = threading.RLock() # RLock for reentrant access
# Reentrant: _prefetch_content releases and re-acquires it around the
# slow fetch while a caller may already hold it.
self._buffer_lock = threading.RLock()
# Plugin rotation state
# Plugin rotation, and the position of the next plugin to fetch in it.
self._ordered_plugins: List[str] = []
self._current_index: int = 0
self._prefetch_index: int = 0
# Update tracking
@@ -97,7 +91,6 @@ class StreamManager:
self.stats = {
'segments_fetched': 0,
'segments_served': 0,
'buffer_swaps': 0,
'fetch_errors': 0,
}
@@ -167,9 +160,7 @@ class StreamManager:
with self._buffer_lock:
return {
'active_count': len(self._active_buffer),
'staging_count': len(self._staging_buffer),
'total_plugins': len(self._ordered_plugins),
'current_index': self._current_index,
'prefetch_index': self._prefetch_index,
'stats': self.stats.copy(),
}
@@ -190,8 +181,9 @@ class StreamManager:
"""
Mark a plugin as having updated data.
Called when a plugin's data changes. Triggers content refresh
for that plugin in the staging buffer.
Only records it in ``_pending_updates``. The refetch happens later,
in :meth:`process_updates` (swap mode) or by dropping the plugin's
caches in :meth:`invalidate_pending_updates` (continuous mode).
Args:
plugin_id: Plugin that was updated
@@ -242,11 +234,6 @@ class StreamManager:
)
return updated
def has_pending_updates(self) -> bool:
"""Check if any plugins have pending updates awaiting processing."""
with self._buffer_lock:
return len(self._pending_updates) > 0
def has_pending_updates_for_visible_segments(self) -> bool:
"""Check if pending updates affect plugins currently in the active buffer."""
with self._buffer_lock:
@@ -309,19 +296,6 @@ class StreamManager:
logger.debug("Processed in-place updates for %d plugins", len(updated_plugins))
def swap_buffers(self) -> None:
"""
Swap active and staging buffers.
Called when staging buffer has updated content ready.
"""
with self._buffer_lock:
if self._staging_buffer:
# True swap: staging becomes active, old active is discarded
self._active_buffer, self._staging_buffer = self._staging_buffer, deque()
self.stats['buffer_swaps'] += 1
logger.debug("Swapped buffers, active now has %d segments", len(self._active_buffer))
def refresh(self) -> None:
"""
Refresh the plugin list and content.
@@ -337,61 +311,49 @@ class StreamManager:
self._refresh_plugin_list()
if len(self._ordered_plugins) != old_count:
logger.info(
logger.debug(
"Plugin list refreshed: %d -> %d plugins",
old_count, len(self._ordered_plugins)
)
def _refresh_plugin_list(self) -> None:
"""Refresh the ordered list of plugins from plugin manager."""
logger.info("=" * 60)
logger.info("REFRESHING PLUGIN LIST FOR VEGAS SCROLL")
logger.info("=" * 60)
"""Refresh the ordered list of plugins from plugin manager.
# Get all enabled plugins
Runs at every cycle start and every ``_refresh_interval`` seconds, so
it logs one INFO summary; the per-plugin decisions are at DEBUG.
"""
available_plugins = []
loaded = 0
if hasattr(self.plugin_manager, 'plugins'):
logger.info(
"Checking %d loaded plugins for Vegas scroll",
len(self.plugin_manager.plugins)
)
loaded = len(self.plugin_manager.plugins)
for plugin_id, plugin in self.plugin_manager.plugins.items():
has_enabled = hasattr(plugin, 'enabled')
is_enabled = getattr(plugin, 'enabled', False)
logger.info(
"[%s] class=%s, has_enabled=%s, enabled=%s",
plugin_id, plugin.__class__.__name__, has_enabled, is_enabled
)
if has_enabled and is_enabled:
# Check vegas content type - skip 'none' unless in STATIC mode
content_type = self.plugin_adapter.get_content_type(plugin, plugin_id)
if not getattr(plugin, 'enabled', False):
logger.debug("[%s] Vegas: skipped (not enabled)", plugin_id)
continue
# Also check display mode - STATIC plugins should be included
# even if their content_type is 'none'
display_mode = VegasDisplayMode.FIXED_SEGMENT
try:
display_mode = plugin.get_vegas_display_mode()
except Exception:
# Plugin error should not abort refresh; use default mode
logger.exception(
"[%s] (%s) get_vegas_display_mode() failed, using default",
plugin_id, plugin.__class__.__name__
)
logger.info(
"[%s] content_type=%s, display_mode=%s",
plugin_id, content_type, display_mode.value
# Content type 'none' is left out, except for STATIC plugins,
# which pause the scroll rather than contributing to it.
content_type = self.plugin_adapter.get_content_type(plugin, plugin_id)
display_mode = VegasDisplayMode.FIXED_SEGMENT
try:
display_mode = plugin.get_vegas_display_mode()
except Exception:
# Plugin error should not abort refresh; use default mode
logger.exception(
"[%s] (%s) get_vegas_display_mode() failed, using default",
plugin_id, plugin.__class__.__name__
)
if content_type != 'none' or display_mode == VegasDisplayMode.STATIC:
available_plugins.append(plugin_id)
logger.info("[%s] --> INCLUDED in Vegas scroll", plugin_id)
else:
logger.info("[%s] --> EXCLUDED from Vegas scroll", plugin_id)
else:
logger.info("[%s] --> SKIPPED (not enabled)", plugin_id)
included = (content_type != 'none'
or display_mode == VegasDisplayMode.STATIC)
logger.debug(
"[%s] Vegas: %s (content_type=%s, display_mode=%s)",
plugin_id, "included" if included else "excluded",
content_type, display_mode.value
)
if included:
available_plugins.append(plugin_id)
else:
logger.warning(
"plugin_manager does not have plugins attribute: %s",
@@ -400,24 +362,20 @@ class StreamManager:
# Apply ordering from config (outside lock for potentially slow operation)
ordered_plugins = self.config.get_ordered_plugins(available_plugins)
logger.info(
"Vegas scroll plugin list: %d available -> %d ordered",
len(available_plugins), len(ordered_plugins)
)
logger.info("Ordered plugins: %s", ordered_plugins)
ordered_plugins = self._apply_priority_weights(ordered_plugins)
# Atomically update shared state under lock to avoid races with prefetchers
with self._buffer_lock:
self._ordered_plugins = ordered_plugins
# Reset indices if needed
if self._current_index >= len(self._ordered_plugins):
self._current_index = 0
if self._prefetch_index >= len(self._ordered_plugins):
self._prefetch_index = 0
logger.info("=" * 60)
slots = (f", {len(ordered_plugins)} slots"
if len(ordered_plugins) != len(set(ordered_plugins)) else "")
logger.info(
"Vegas rotation: %d of %d loaded plugin(s)%s: %s",
len(set(ordered_plugins)), loaded, slots, ', '.join(ordered_plugins)
)
def _plugin_weight(self, plugin_id: str) -> int:
"""Slots per cycle for one plugin.
@@ -490,7 +448,7 @@ class StreamManager:
schedule = self._unclump_seam(schedule)
boosted = {p: w for p, w in weights.items() if w > 1}
logger.info(
logger.debug(
"Vegas rotation weighted: %d slots for %d plugins (boosted: %s)",
len(schedule), len(ordered), boosted)
return schedule
@@ -607,11 +565,6 @@ class StreamManager:
ContentSegment or None if fetch failed
"""
try:
logger.info("=" * 60)
logger.info("[%s] FETCHING CONTENT", plugin_id)
logger.info("=" * 60)
# Get plugin instance
if not hasattr(self.plugin_manager, 'plugins'):
logger.warning("[%s] plugin_manager has no plugins attribute", plugin_id)
return None
@@ -621,18 +574,11 @@ class StreamManager:
logger.warning("[%s] Plugin not found in plugin_manager.plugins", plugin_id)
return None
logger.info(
"[%s] Plugin found: class=%s, enabled=%s",
plugin_id, plugin.__class__.__name__, getattr(plugin, 'enabled', 'N/A')
)
# Get display mode from plugin
display_mode = VegasDisplayMode.FIXED_SEGMENT
try:
display_mode = plugin.get_vegas_display_mode()
logger.info("[%s] Display mode: %s", plugin_id, display_mode.value)
except (AttributeError, TypeError) as e:
logger.info(
logger.debug(
"[%s] get_vegas_display_mode() not available: %s (using FIXED_SEGMENT)",
plugin_id, e
)
@@ -644,21 +590,21 @@ class StreamManager:
segment = ContentSegment(
plugin_id=plugin_id,
images=[], # No images needed for static pause
total_width=0,
display_mode=display_mode
)
self.stats['segments_fetched'] += 1
logger.info(
logger.debug(
"[%s] Created STATIC placeholder (pause trigger)",
plugin_id
)
return segment
# Get content via adapter for SCROLL/FIXED_SEGMENT modes
logger.info("[%s] Calling plugin_adapter.get_content()...", plugin_id)
images = self.plugin_adapter.get_content(plugin, plugin_id)
if not images:
logger.warning("[%s] NO CONTENT RETURNED from plugin_adapter", plugin_id)
# The adapter already warns when every content path failed;
# an empty result is otherwise routine (nothing scheduled).
logger.debug("[%s] No Vegas content this cycle", plugin_id)
return None
# Calculate total width
@@ -667,17 +613,14 @@ class StreamManager:
segment = ContentSegment(
plugin_id=plugin_id,
images=images,
total_width=total_width,
display_mode=display_mode
)
self.stats['segments_fetched'] += 1
logger.info(
"[%s] SEGMENT CREATED: %d images, %dpx total, mode=%s",
logger.debug(
"[%s] Segment: %d image(s), %dpx, mode=%s",
plugin_id, len(images), total_width, display_mode.value
)
logger.info("=" * 60)
return segment
except Exception:
@@ -696,24 +639,6 @@ class StreamManager:
if len(self._active_buffer) < low_water:
self._prefetch_content(count=low_water - len(self._active_buffer))
def get_all_content_for_composition(self) -> List[Image.Image]:
"""
Get all buffered content as a flat list of images.
Skips STATIC segments as they don't have images to compose.
Prefer get_grouped_content_for_composition(): flattening loses the
plugin boundaries, which is what tells the compositor where a
separator belongs and where it does not.
Returns:
List of all images in buffer order
"""
all_images = []
for _plugin_id, images in self.get_grouped_content_for_composition():
all_images.extend(images)
return all_images
def get_grouped_content_for_composition(self) -> List[Tuple[str, List[Image.Image]]]:
"""
Get buffered content grouped by the plugin that produced it.
@@ -826,8 +751,6 @@ class StreamManager:
"""Reset the stream manager state."""
with self._buffer_lock:
self._active_buffer.clear()
self._staging_buffer.clear()
self._current_index = 0
self._prefetch_index = 0
self._pending_updates.clear()