mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-10 17:16:36 +00:00
* fix(display): tear down Vegas mode on controller cleanup DisplayController.cleanup() never called VegasModeCoordinator.cleanup(), so the Vegas teardown (stop, pipeline/stream reset, adapter cache drop) was unreachable. Call it before the display manager is cleaned up, and skip it when Vegas was never created. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(vegas): default max_cycle_duration to the documented 240s The template, the web UI help, CONFIG_REFERENCE and the controller all say 240, but the code defaulted to 600 in two places, so a config without the key ran Vegas iterations 2.5x longer than documented. from_config now falls back to the dataclass field defaults instead of repeating each one, so the two copies can no longer drift, and the controller's follower scroll-speed default reads VegasModeConfig's. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(display): let run.py -d show display_manager's DEBUG output display_manager pinned its logger to INFO at import, overriding the root level, so debug mode never showed its DEBUG lines. Use get_logger() from src.logging_config like the rest of the core and leave the level to the logging setup. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(display): run each startup validation check once StartupValidator.validate_all() ran twice at boot, before and after the plugin manager was created, so every config, cache, display and systemd-unit warning was logged twice. The second pass now runs only the plugin checks. Drop the commented-out raise_on_errors line. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(vegas): one INFO line per plugin-list refresh StreamManager logged "=" * 60 banners and a line per plugin (INCLUDED, SKIPPED, FETCHING CONTENT, SEGMENT CREATED) at INFO on every refresh and fetch, i.e. at each cycle start and every 30s. Log one INFO summary of the rotation per refresh and move the per-plugin detail, the weighting breakdown and "no content this cycle" to DEBUG (the adapter still warns when every content path fails). Also drop the check/cross marks from the controller's log messages. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(vegas): drop the per-iteration static-mode plugin scan run_iteration() rebuilt _static_mode_plugins on every iteration, asking every plugin for its display mode and logging the set at INFO, but nothing ever read it: static pauses are triggered by _check_static_plugin_trigger() from the next segment. Delete it, the coordinator's get_ordered_plugins() that only it used, and the write-only _static_pause_plugin / _static_pause_start. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(vegas): remove the staging buffer that was never filled StreamManager and RenderPipeline carried a double-buffer design that nothing used: _staging_buffer was only ever cleared or swapped, so swap_buffers() never did anything and should_recompose()'s staging_count > 0 branch was dead, and _active_scroll_image, _staging_scroll_image, _is_rendering, _last_frame_time and _frame_interval were written but never read. Delete the machinery and rewrite the docstrings around what actually carries updates: _pending_updates, consumed by process_updates() in swap mode and invalidate_pending_updates() in continuous mode. should_recompose() no longer builds a buffer-status dict every frame. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(display): tidy the display controller without changing behaviour - Import VegasModeCoordinator locally instead of through module globals (there is no circular import to avoid). - Drop hasattr() checks on attributes PluginManager.__init__ always sets (plugin_executor, plugin_last_update, get_plugin_lock, run_scheduled_updates*, stop_update_worker) and the dead "older manager" fallbacks; keep the health_tracker None checks, now via _health_tracker(). - Extract _display_once() for the per-frame display call both render loops copied, _advance_on_demand() for the two on-demand rotations, _reset_on_demand_fields() for the error and clear paths, and _timezone() / _in_window() for the two schedule checks. - Remove always-true conditions and the unreachable non-plugin else branch in run(), and read _was_display_active / _last_published_mode / vegas_coordinator directly now that __init__ declares them. - Declare the follower render state in __init__, name its tuning constants, add _follower_sign(), and share the 90/s sync send interval with the render pipeline (SYNC_SEND_INTERVAL). - Delete history narration and the "Opt #N" labels, fix the comment that called _scroll_speed constant (hot reload updates it), and drop a startup timing log that measured nothing. - render_pipeline / plugin_adapter: read display_manager.width/height as the properties they are, drop an empty TYPE_CHECKING block, an aliased threading import and a duplicated `if result and self.sync_manager:`. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(display): trim dead code from display_manager - Add _new_canvas() for the image/draw/fontmode="1" setup that was copied six times. - Call resolve_double_sided() and compose_pixel_mapper_config() directly instead of through a module alias and a passthrough method, and replace the comment that said the passthrough read class attributes. - Delete the unused _initialized flag and _ORIENTATION_ROTATE_DEGREES alias (no core or monorepo reader; tests stop resetting the flag), the test pattern's unreachable no-matrix branch (it only runs once the matrix exists), `del old_image # help GC` (a no-op on a local), a duplicated early return in process_deferred_updates, and stale comments. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(vegas): remove unread fields and test-only helpers, fix docstrings - ContentSegment: drop total_width, fetched_at, is_stale, image_count and is_static, none of which is read. - StreamManager: drop _current_index (never advanced) and the test-only get_all_content_for_composition() and has_pending_updates(); VegasModeConfig: drop the test-only is_plugin_included(). - geometry.find_blank_cut() has had no production caller since the crop moved to item boundaries; delete it and its tests. - PluginAdapter: the _finalize docstring described separator_width between every image, and _crop_to_budget's said cuts snap to the nearest blank column; both now describe what the code does. - Coordinator: the static-pause interrupt log no longer blames follower mode for every interrupt, and set_update_callback names the callback the controller actually wires. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(scroll): correct ScrollHelper comments and drop dead branches - Four comments said the strip always starts with display_width of blank; it does only when lead_gap is None (Vegas passes its own). - Delete the "Width calculation mismatch" warning: the image is created at the calculated width, so the two can never differ. - Remove the two scroll_delay <= 0 fallbacks (which disagreed with each other): set_scroll_delay clamps it to at least 0.001 and nothing in core or the plugin monorepo assigns it directly. - Trim the scipy history from the blend docstring. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * refactor(run): drop a redundant comment Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(display): log set_scrolling_state only when it changes Vegas and scrolling plugins set the scrolling state every frame, so once display_manager's DEBUG output became visible in debug mode it printed "Scrolling state set to: True" about 120 times a second. Log only when the value differs from the previous one; the state, activity timestamp and frame hold still update on every call. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * docs(changelog): display-vegas Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
755 lines
29 KiB
Python
755 lines
29 KiB
Python
"""
|
|
Stream Manager for Vegas Mode
|
|
|
|
Manages plugin content streaming with look-ahead buffering. Maintains a queue
|
|
of plugin content that's ready to be rendered, prefetching 1-2 plugins ahead
|
|
of the current scroll position.
|
|
|
|
Supports three display modes:
|
|
- SCROLL: Continuous scrolling content
|
|
- FIXED_SEGMENT: Fixed block that scrolls by
|
|
- STATIC: Pause scroll to display (marked for coordinator handling)
|
|
"""
|
|
|
|
import logging
|
|
import threading
|
|
import time
|
|
from typing import Optional, List, Dict, Any, Deque, Tuple, TYPE_CHECKING
|
|
from collections import deque
|
|
from dataclasses import dataclass, field
|
|
from PIL import Image
|
|
|
|
from src.vegas_mode.config import VegasModeConfig
|
|
from src.vegas_mode.plugin_adapter import PluginAdapter
|
|
from src.plugin_system.base_plugin import VegasDisplayMode
|
|
|
|
if TYPE_CHECKING:
|
|
from src.plugin_system.plugin_manager import PluginManager
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
@dataclass
|
|
class ContentSegment:
|
|
"""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]
|
|
display_mode: VegasDisplayMode = field(default=VegasDisplayMode.FIXED_SEGMENT)
|
|
|
|
|
|
class StreamManager:
|
|
"""
|
|
Manages streaming of plugin content for Vegas scroll mode.
|
|
|
|
Key responsibilities:
|
|
- 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__(
|
|
self,
|
|
config: VegasModeConfig,
|
|
plugin_manager: 'PluginManager',
|
|
plugin_adapter: PluginAdapter
|
|
):
|
|
"""
|
|
Initialize the stream manager.
|
|
|
|
Args:
|
|
config: Vegas mode configuration
|
|
plugin_manager: Plugin manager for accessing plugins
|
|
plugin_adapter: Adapter for getting plugin content
|
|
"""
|
|
self.config = config
|
|
self.plugin_manager = plugin_manager
|
|
self.plugin_adapter = plugin_adapter
|
|
|
|
# Segments composed into the current cycle (swap mode only).
|
|
self._active_buffer: Deque[ContentSegment] = deque()
|
|
# 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, and the position of the next plugin to fetch in it.
|
|
self._ordered_plugins: List[str] = []
|
|
self._prefetch_index: int = 0
|
|
|
|
# Update tracking
|
|
self._pending_updates: Dict[str, bool] = {}
|
|
self._last_refresh: float = 0.0
|
|
self._refresh_interval: float = 30.0 # Refresh plugin list every 30s
|
|
|
|
# Statistics
|
|
self.stats = {
|
|
'segments_fetched': 0,
|
|
'segments_served': 0,
|
|
'fetch_errors': 0,
|
|
}
|
|
|
|
logger.info("StreamManager initialized with buffer_ahead=%d", config.buffer_ahead)
|
|
|
|
def initialize(self) -> bool:
|
|
"""
|
|
Initialize the stream manager with current plugin list.
|
|
|
|
Returns:
|
|
True if initialized successfully with at least one plugin
|
|
"""
|
|
self._refresh_plugin_list()
|
|
|
|
if not self._ordered_plugins:
|
|
logger.warning("No plugins available for Vegas scroll")
|
|
return False
|
|
|
|
# Fill the buffer to a whole cycle's worth of plugins. This used to be
|
|
# buffer_ahead + 1, which conflated prefetch depth with cycle size and
|
|
# meant a 20-plugin install only showed 3 plugins before recomposing.
|
|
self._prefetch_content(
|
|
count=min(self.config.plugins_per_cycle, len(self._ordered_plugins)))
|
|
|
|
logger.info(
|
|
"StreamManager initialized with %d plugins, %d segments buffered",
|
|
len(self._ordered_plugins), len(self._active_buffer)
|
|
)
|
|
return len(self._active_buffer) > 0
|
|
|
|
def get_next_segment(self) -> Optional[ContentSegment]:
|
|
"""
|
|
Get the next content segment for rendering.
|
|
|
|
Returns:
|
|
ContentSegment or None if buffer is empty
|
|
"""
|
|
with self._buffer_lock:
|
|
if not self._active_buffer:
|
|
# Try to fetch more content
|
|
self._prefetch_content(count=1)
|
|
if not self._active_buffer:
|
|
return None
|
|
|
|
segment = self._active_buffer.popleft()
|
|
self.stats['segments_served'] += 1
|
|
|
|
# Trigger prefetch to maintain buffer
|
|
self._ensure_buffer_filled()
|
|
|
|
return segment
|
|
|
|
def peek_next_segment(self) -> Optional[ContentSegment]:
|
|
"""
|
|
Peek at the next segment without removing it.
|
|
|
|
Returns:
|
|
ContentSegment or None if buffer is empty
|
|
"""
|
|
with self._buffer_lock:
|
|
if self._active_buffer:
|
|
return self._active_buffer[0]
|
|
return None
|
|
|
|
def get_buffer_status(self) -> Dict[str, Any]:
|
|
"""Get current buffer status for monitoring."""
|
|
with self._buffer_lock:
|
|
return {
|
|
'active_count': len(self._active_buffer),
|
|
'total_plugins': len(self._ordered_plugins),
|
|
'prefetch_index': self._prefetch_index,
|
|
'stats': self.stats.copy(),
|
|
}
|
|
|
|
def get_active_plugin_ids(self) -> List[str]:
|
|
"""
|
|
Get list of plugin IDs currently in the active buffer.
|
|
|
|
Thread-safe accessor for render pipeline.
|
|
|
|
Returns:
|
|
List of plugin IDs in buffer order
|
|
"""
|
|
with self._buffer_lock:
|
|
return [seg.plugin_id for seg in self._active_buffer]
|
|
|
|
def mark_plugin_updated(self, plugin_id: str) -> None:
|
|
"""
|
|
Mark a plugin as having updated data.
|
|
|
|
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
|
|
"""
|
|
with self._buffer_lock:
|
|
self._pending_updates[plugin_id] = True
|
|
|
|
logger.debug("Plugin %s marked for update", plugin_id)
|
|
|
|
def invalidate_pending_updates(self) -> List[str]:
|
|
"""
|
|
Drop cached content for plugins whose data changed, without refetching.
|
|
|
|
The continuous-scroll counterpart to :meth:`process_updates`. That method
|
|
belongs to the swap path: it refetches immediately and merges into the
|
|
active buffer, which continuous mode bypasses entirely, and doing that
|
|
work on the render thread would hitch the scroll.
|
|
|
|
Here it is enough to clear the caches and let the plugin come round in
|
|
the rotation, which recomposes it from current data a moment later. Left
|
|
uncalled, ``_pending_updates`` simply accumulates and no visual ever
|
|
refreshes — a game that was live last night keeps being drawn as live.
|
|
|
|
Returns:
|
|
The plugin ids whose caches were dropped.
|
|
"""
|
|
with self._buffer_lock:
|
|
if not self._pending_updates:
|
|
return []
|
|
updated = list(self._pending_updates.keys())
|
|
self._pending_updates.clear()
|
|
|
|
plugins = getattr(self.plugin_manager, 'plugins', {})
|
|
for plugin_id in updated:
|
|
try:
|
|
self.plugin_adapter.invalidate_cache(plugin_id)
|
|
plugin = plugins.get(plugin_id)
|
|
if plugin is not None:
|
|
self.plugin_adapter.invalidate_plugin_scroll_cache(
|
|
plugin, plugin_id)
|
|
except Exception: # pylint: disable=broad-except
|
|
logger.exception(
|
|
"[%s] Could not invalidate cached content", plugin_id)
|
|
|
|
logger.info(
|
|
"Vegas: dropped cached content for %d updated plugin(s): %s",
|
|
len(updated), ', '.join(updated)
|
|
)
|
|
return updated
|
|
|
|
def has_pending_updates_for_visible_segments(self) -> bool:
|
|
"""Check if pending updates affect plugins currently in the active buffer."""
|
|
with self._buffer_lock:
|
|
if not self._pending_updates:
|
|
return False
|
|
active_ids = {
|
|
seg.plugin_id for seg in self._active_buffer if seg.images
|
|
}
|
|
return bool(active_ids & self._pending_updates.keys())
|
|
|
|
def process_updates(self) -> None:
|
|
"""
|
|
Process pending plugin updates.
|
|
|
|
Performs in-place update of segments in the active buffer,
|
|
preserving non-updated plugins and their order.
|
|
"""
|
|
with self._buffer_lock:
|
|
if not self._pending_updates:
|
|
return
|
|
|
|
updated_plugins = list(self._pending_updates.keys())
|
|
self._pending_updates.clear()
|
|
|
|
# Fetch fresh content for each updated plugin (outside lock for slow ops)
|
|
refreshed_segments = {}
|
|
for plugin_id in updated_plugins:
|
|
self.plugin_adapter.invalidate_cache(plugin_id)
|
|
|
|
# Clear the plugin's scroll_helper cache so the visual is rebuilt
|
|
# from fresh data (affects stocks, news, odds-ticker, etc.)
|
|
plugin = None
|
|
if hasattr(self.plugin_manager, 'plugins'):
|
|
plugin = self.plugin_manager.plugins.get(plugin_id)
|
|
if plugin:
|
|
self.plugin_adapter.invalidate_plugin_scroll_cache(plugin, plugin_id)
|
|
|
|
segment = self._fetch_plugin_content(plugin_id)
|
|
if segment:
|
|
refreshed_segments[plugin_id] = segment
|
|
|
|
# In-place merge: replace segments in active buffer
|
|
with self._buffer_lock:
|
|
# Build new buffer preserving order, replacing updated segments
|
|
new_buffer: Deque[ContentSegment] = deque()
|
|
seen_plugins: set = set()
|
|
|
|
for segment in self._active_buffer:
|
|
if segment.plugin_id in refreshed_segments:
|
|
# Replace with refreshed segment (only once per plugin)
|
|
if segment.plugin_id not in seen_plugins:
|
|
new_buffer.append(refreshed_segments[segment.plugin_id])
|
|
seen_plugins.add(segment.plugin_id)
|
|
# Skip duplicate entries for same plugin
|
|
else:
|
|
# Keep non-updated segment
|
|
new_buffer.append(segment)
|
|
|
|
self._active_buffer = new_buffer
|
|
|
|
logger.debug("Processed in-place updates for %d plugins", len(updated_plugins))
|
|
|
|
def refresh(self) -> None:
|
|
"""
|
|
Refresh the plugin list and content.
|
|
|
|
Called periodically to pick up new plugins or config changes.
|
|
"""
|
|
current_time = time.time()
|
|
if current_time - self._last_refresh < self._refresh_interval:
|
|
return
|
|
|
|
self._last_refresh = current_time
|
|
old_count = len(self._ordered_plugins)
|
|
self._refresh_plugin_list()
|
|
|
|
if len(self._ordered_plugins) != old_count:
|
|
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.
|
|
|
|
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'):
|
|
loaded = len(self.plugin_manager.plugins)
|
|
for plugin_id, plugin in self.plugin_manager.plugins.items():
|
|
if not getattr(plugin, 'enabled', False):
|
|
logger.debug("[%s] Vegas: skipped (not enabled)", plugin_id)
|
|
continue
|
|
|
|
# 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__
|
|
)
|
|
|
|
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",
|
|
type(self.plugin_manager).__name__
|
|
)
|
|
|
|
# Apply ordering from config (outside lock for potentially slow operation)
|
|
ordered_plugins = self.config.get_ordered_plugins(available_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
|
|
if self._prefetch_index >= len(self._ordered_plugins):
|
|
self._prefetch_index = 0
|
|
|
|
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.
|
|
|
|
A plugin may answer for itself via get_vegas_priority_weight() -- the
|
|
only way favorite-team awareness can reach here, since the core can see
|
|
that a game is live but not whose. When it declines (returns None, the
|
|
default), live content earns ``live_weight`` and everything else 1.
|
|
"""
|
|
plugin = None
|
|
try:
|
|
plugin = self.plugin_manager.plugins.get(plugin_id)
|
|
except (AttributeError, TypeError):
|
|
return 1
|
|
if plugin is None:
|
|
return 1
|
|
|
|
try:
|
|
if hasattr(plugin, 'get_vegas_priority_weight'):
|
|
declared = plugin.get_vegas_priority_weight()
|
|
if declared is not None:
|
|
return max(1, min(10, int(declared)))
|
|
except Exception:
|
|
# Deliberately falls through to the core's own live check rather
|
|
# than demoting to 1. The plugin's weight calculation is broken,
|
|
# but has_live_priority() and has_live_content() are separate
|
|
# methods guarded separately below -- a plugin that genuinely has
|
|
# a live game should still get live_weight for it.
|
|
logger.exception("[%s] get_vegas_priority_weight() failed", plugin_id)
|
|
|
|
try:
|
|
if (hasattr(plugin, 'has_live_priority')
|
|
and hasattr(plugin, 'has_live_content')
|
|
and plugin.has_live_priority()
|
|
and plugin.has_live_content()):
|
|
return self.config.live_weight
|
|
except Exception:
|
|
logger.exception("[%s] live-content check failed", plugin_id)
|
|
return 1
|
|
|
|
def _apply_priority_weights(self, ordered: List[str]) -> List[str]:
|
|
"""Expand the rotation so weighted plugins take several turns per cycle.
|
|
|
|
Smooth Weighted Round-Robin, the same scheduler the sports plugins use
|
|
to rotate their own games: a plugin of weight N appears N times per
|
|
cycle, and the repeats are spaced through the cycle rather than
|
|
clumped, so a live score is never three-in-a-row followed by a long
|
|
silence.
|
|
|
|
Returns the input unchanged when nothing is weighted, which is both the
|
|
common case and the pre-existing behaviour.
|
|
"""
|
|
if not ordered or not self.config.live_in_ticker:
|
|
return ordered
|
|
|
|
weights = {pid: self._plugin_weight(pid) for pid in ordered}
|
|
total = sum(weights.values())
|
|
if total <= len(ordered):
|
|
return ordered # nothing boosted; plain round robin
|
|
|
|
current = {pid: 0 for pid in ordered}
|
|
schedule: List[str] = []
|
|
for _ in range(total):
|
|
for pid in ordered:
|
|
current[pid] += weights[pid]
|
|
picked = max(current, key=lambda p: current[p])
|
|
current[picked] -= total
|
|
schedule.append(picked)
|
|
|
|
schedule = self._unclump_seam(schedule)
|
|
|
|
boosted = {p: w for p, w in weights.items() if w > 1}
|
|
logger.debug(
|
|
"Vegas rotation weighted: %d slots for %d plugins (boosted: %s)",
|
|
len(schedule), len(ordered), boosted)
|
|
return schedule
|
|
|
|
@staticmethod
|
|
def _unclump_seam(schedule: List[str]) -> List[str]:
|
|
"""Stop the heaviest plugin sitting on both ends of the cycle.
|
|
|
|
Smooth Weighted Round-Robin spaces repeats well *within* a pass, but
|
|
it schedules the heaviest item first and often last too. The strip
|
|
loops, so those two are neighbours: the one place the marquee shows
|
|
the same plugin twice running is the seam between cycles.
|
|
|
|
Rotating the list cannot fix this. Rotation preserves the cyclic order
|
|
exactly, so it only moves where the seam is drawn, not the adjacency
|
|
itself. The trailing entry has to be swapped with one from the middle
|
|
whose neighbours differ from it, which breaks the pair without
|
|
creating another.
|
|
|
|
Left alone when no such position exists -- a rotation short enough or
|
|
lopsided enough to have none is one where the plugin is unavoidably
|
|
adjacent to itself anyway.
|
|
"""
|
|
if len(schedule) < 3 or schedule[0] != schedule[-1]:
|
|
return schedule
|
|
|
|
repeated = schedule[-1]
|
|
size = len(schedule)
|
|
|
|
def cyclic_doubles(seq) -> int:
|
|
return sum(1 for i in range(size) if seq[i] == seq[(i + 1) % size])
|
|
|
|
def clearance(seq, value) -> int:
|
|
"""Smallest cyclic gap between appearances of `value`."""
|
|
at = [i for i, v in enumerate(seq) if v == value]
|
|
if len(at) < 2:
|
|
return size
|
|
return min(min((b - a) % size, (a - b) % size)
|
|
for i, a in enumerate(at) for b in at[i + 1:])
|
|
|
|
# Try each swap and judge the result, rather than reasoning about which
|
|
# neighbours the two moved elements will end up with. That reasoning is
|
|
# where the first version went wrong: it guarded the slot `repeated`
|
|
# moves into but not the one the displaced element lands in, so
|
|
# ['a','b','c','d','x','y','x','a'] came back ending ['x','x'] -- the
|
|
# seam duplicate traded for a fresh one.
|
|
best = None
|
|
best_clearance = -1
|
|
for j in range(1, size - 1):
|
|
candidate = list(schedule)
|
|
candidate[j], candidate[-1] = candidate[-1], candidate[j]
|
|
if cyclic_doubles(candidate):
|
|
continue
|
|
# Among the repairs that work, prefer the one that leaves the
|
|
# boosted plugin most evenly spread; taking the first that merely
|
|
# fits moved a repeat from a gap of 7 into a gap of 2.
|
|
spread = clearance(candidate, repeated)
|
|
if spread > best_clearance:
|
|
best, best_clearance = candidate, spread
|
|
|
|
# None exists when the value is unavoidably adjacent to itself -- a
|
|
# plugin holding most of the slots has to be. Schedule it as it is
|
|
# rather than refuse.
|
|
return best if best is not None else schedule
|
|
|
|
def _prefetch_content(self, count: int = 1) -> None:
|
|
"""
|
|
Prefetch content for upcoming plugins.
|
|
|
|
Args:
|
|
count: Number of plugins to prefetch
|
|
"""
|
|
with self._buffer_lock:
|
|
if not self._ordered_plugins:
|
|
return
|
|
|
|
for _ in range(count):
|
|
if len(self._active_buffer) >= self.config.plugins_per_cycle:
|
|
break
|
|
|
|
# Ensure index is valid (guard against empty list)
|
|
num_plugins = len(self._ordered_plugins)
|
|
if num_plugins == 0:
|
|
break
|
|
|
|
plugin_id = self._ordered_plugins[self._prefetch_index]
|
|
|
|
# Release lock for potentially slow content fetch
|
|
self._buffer_lock.release()
|
|
try:
|
|
segment = self._fetch_plugin_content(plugin_id)
|
|
finally:
|
|
self._buffer_lock.acquire()
|
|
|
|
if segment:
|
|
self._active_buffer.append(segment)
|
|
|
|
# Revalidate num_plugins after reacquiring lock (may have changed)
|
|
num_plugins = len(self._ordered_plugins)
|
|
if num_plugins == 0:
|
|
break
|
|
|
|
# Advance prefetch index (thread-safe within lock)
|
|
self._prefetch_index = (self._prefetch_index + 1) % num_plugins
|
|
|
|
def _fetch_plugin_content(self, plugin_id: str) -> Optional[ContentSegment]:
|
|
"""
|
|
Fetch content from a specific plugin.
|
|
|
|
Args:
|
|
plugin_id: Plugin to fetch from
|
|
|
|
Returns:
|
|
ContentSegment or None if fetch failed
|
|
"""
|
|
try:
|
|
if not hasattr(self.plugin_manager, 'plugins'):
|
|
logger.warning("[%s] plugin_manager has no plugins attribute", plugin_id)
|
|
return None
|
|
|
|
plugin = self.plugin_manager.plugins.get(plugin_id)
|
|
if not plugin:
|
|
logger.warning("[%s] Plugin not found in plugin_manager.plugins", plugin_id)
|
|
return None
|
|
|
|
display_mode = VegasDisplayMode.FIXED_SEGMENT
|
|
try:
|
|
display_mode = plugin.get_vegas_display_mode()
|
|
except (AttributeError, TypeError) as e:
|
|
logger.debug(
|
|
"[%s] get_vegas_display_mode() not available: %s (using FIXED_SEGMENT)",
|
|
plugin_id, e
|
|
)
|
|
|
|
# For STATIC mode, we create a placeholder segment
|
|
# The actual content will be displayed by coordinator during pause
|
|
if display_mode == VegasDisplayMode.STATIC:
|
|
# Create minimal placeholder - coordinator handles actual display
|
|
segment = ContentSegment(
|
|
plugin_id=plugin_id,
|
|
images=[], # No images needed for static pause
|
|
display_mode=display_mode
|
|
)
|
|
self.stats['segments_fetched'] += 1
|
|
logger.debug(
|
|
"[%s] Created STATIC placeholder (pause trigger)",
|
|
plugin_id
|
|
)
|
|
return segment
|
|
|
|
# Get content via adapter for SCROLL/FIXED_SEGMENT modes
|
|
images = self.plugin_adapter.get_content(plugin, plugin_id)
|
|
if not images:
|
|
# 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
|
|
total_width = sum(img.width for img in images)
|
|
|
|
segment = ContentSegment(
|
|
plugin_id=plugin_id,
|
|
images=images,
|
|
display_mode=display_mode
|
|
)
|
|
|
|
self.stats['segments_fetched'] += 1
|
|
logger.debug(
|
|
"[%s] Segment: %d image(s), %dpx, mode=%s",
|
|
plugin_id, len(images), total_width, display_mode.value
|
|
)
|
|
return segment
|
|
|
|
except Exception:
|
|
logger.exception("[%s] ERROR fetching content", plugin_id)
|
|
self.stats['fetch_errors'] += 1
|
|
return None
|
|
|
|
def _ensure_buffer_filled(self) -> None:
|
|
"""
|
|
Top the buffer back up after segments have been served.
|
|
|
|
buffer_ahead is the low-water mark only; plugins_per_cycle is the
|
|
ceiling and is enforced inside _prefetch_content.
|
|
"""
|
|
low_water = min(self.config.buffer_ahead, self.config.plugins_per_cycle)
|
|
if len(self._active_buffer) < low_water:
|
|
self._prefetch_content(count=low_water - len(self._active_buffer))
|
|
|
|
def get_grouped_content_for_composition(self) -> List[Tuple[str, List[Image.Image]]]:
|
|
"""
|
|
Get buffered content grouped by the plugin that produced it.
|
|
|
|
The grouping matters: separator_width is meant to mark the handoff from
|
|
one plugin to the next, not to sit between every row a single plugin
|
|
contributes. A per-row ticker like the F1 scoreboard returns over a
|
|
hundred images that it renders 4px apart internally, so flattening them
|
|
into one list and applying a uniform gap forced 32px between each of
|
|
its rows — both inconsistent with how the plugin looks standalone, and
|
|
a large hidden addition to the width it occupies.
|
|
|
|
Skips STATIC segments, which trigger a pause rather than contributing
|
|
scroll content, and segments left with no images.
|
|
|
|
Returns:
|
|
List of (plugin_id, images) in buffer order
|
|
"""
|
|
grouped: List[Tuple[str, List[Image.Image]]] = []
|
|
with self._buffer_lock:
|
|
for segment in self._active_buffer:
|
|
if segment.display_mode == VegasDisplayMode.STATIC:
|
|
continue
|
|
if not segment.images:
|
|
continue
|
|
grouped.append((segment.plugin_id, list(segment.images)))
|
|
return grouped
|
|
|
|
def take_next_group(
|
|
self, count: Optional[int] = None, offscreen_only: bool = False
|
|
) -> List[Tuple[str, Optional[List[Image.Image]]]]:
|
|
"""
|
|
Fetch and hand over the next slice of the rotation.
|
|
|
|
For continuous scrolling, where the strip is extended rather than
|
|
replaced. Advances the rotation index so plugins come round in order
|
|
across an unbroken strip, and bypasses the active buffer entirely — that
|
|
buffer exists to stage a *replacement* cycle, which continuous mode has
|
|
no use for.
|
|
|
|
Args:
|
|
count: Number of plugins to gather, defaulting to plugins_per_cycle
|
|
offscreen_only: Only use content paths that avoid the shared display
|
|
canvas, for use off the render thread
|
|
|
|
Returns:
|
|
Ordered list of (plugin_id, images). ``images`` is None when the
|
|
plugin could not be served under ``offscreen_only``, so the caller
|
|
can fetch just those on the render thread while keeping the order.
|
|
"""
|
|
if count is None:
|
|
count = self.config.plugins_per_cycle
|
|
|
|
self.refresh()
|
|
|
|
with self._buffer_lock:
|
|
if not self._ordered_plugins:
|
|
return []
|
|
total = len(self._ordered_plugins)
|
|
ids = []
|
|
for _ in range(min(max(1, count), total)):
|
|
ids.append(self._ordered_plugins[self._prefetch_index])
|
|
self._prefetch_index = (self._prefetch_index + 1) % total
|
|
|
|
plugins = getattr(self.plugin_manager, 'plugins', {})
|
|
group: List[Tuple[str, Optional[List[Image.Image]]]] = []
|
|
|
|
for plugin_id in ids:
|
|
plugin = plugins.get(plugin_id)
|
|
if not plugin:
|
|
continue
|
|
try:
|
|
images = self.plugin_adapter.get_content(
|
|
plugin, plugin_id, offscreen_only=offscreen_only)
|
|
except Exception:
|
|
logger.exception("[%s] ERROR fetching content", plugin_id)
|
|
self.stats['fetch_errors'] += 1
|
|
continue
|
|
if images:
|
|
self.stats['segments_fetched'] += 1
|
|
group.append((plugin_id, images if images else None))
|
|
|
|
return group
|
|
|
|
def advance_cycle(self) -> None:
|
|
"""
|
|
Advance to next cycle by clearing the active buffer.
|
|
|
|
Called when a scroll cycle completes to allow fresh content
|
|
to be fetched for the next cycle. Does not reset indices,
|
|
so prefetching continues from the current position in the
|
|
plugin order.
|
|
"""
|
|
with self._buffer_lock:
|
|
consumed_count = len(self._active_buffer)
|
|
self._active_buffer.clear()
|
|
logger.debug("Advanced cycle, cleared %d segments", consumed_count)
|
|
|
|
def reset(self) -> None:
|
|
"""Reset the stream manager state."""
|
|
with self._buffer_lock:
|
|
self._active_buffer.clear()
|
|
self._prefetch_index = 0
|
|
self._pending_updates.clear()
|
|
|
|
self.plugin_adapter.invalidate_cache()
|
|
logger.info("StreamManager reset")
|
|
|
|
def cleanup(self) -> None:
|
|
"""Clean up resources."""
|
|
self.reset()
|
|
self.plugin_adapter.cleanup()
|
|
logger.debug("StreamManager cleanup complete")
|