mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-04 22:35:08 +00:00
feat(vegas): live elements update in place while they scroll
One background worker (src/vegas_mode/live_worker.py) redraws a plugin's live elements when its data epoch moves on (update listener) or on their refresh_hz, nearest the screen first, and hands changed pixels lock-free to the render thread, which copies them into the strip between frames (RenderPipeline.apply_live_patches, ScrollHelper.patch_columns): at most four patches or two screens of bytes a frame, no drawing or locks there. The worker takes over group prefetch once a live element is placed, runs inside the render gate, and is supervised. Update tick 1s while live elements exist. Web UI switch for live_refresh. OFFSCREEN_RENDERING.md describes what was built and why SegmentStrip was not needed. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -780,6 +780,43 @@ class ScrollHelper:
|
||||
)
|
||||
return cut
|
||||
|
||||
def patch_columns(self, x: int, pixels: np.ndarray) -> int:
|
||||
"""Overwrite the strip's columns from ``x`` with ``pixels``, in place.
|
||||
|
||||
What a live Vegas element update is (src/vegas_mode/elements.py): the
|
||||
strip keeps its width, the scroll keeps its position, and only these
|
||||
columns change. Call it between frames on the thread that draws them;
|
||||
every frame copies its slice out of the strip (get_visible_portion),
|
||||
so no frame already handed on can see half a patch.
|
||||
|
||||
Clipped to the strip at both ends. Refused (0) for an array this
|
||||
helper may not write -- the multi-display follower adopts a read-only
|
||||
one -- or for pixels of another height. The PIL image is deferred, so
|
||||
a later read of cached_image shows the patch.
|
||||
|
||||
Args:
|
||||
x: Strip column of the first column of ``pixels``
|
||||
pixels: uint8 array (height, width, 3)
|
||||
|
||||
Returns:
|
||||
Bytes written.
|
||||
"""
|
||||
strip = self.cached_array
|
||||
if strip is None or not strip.flags.writeable:
|
||||
return 0
|
||||
if pixels.ndim != 3 or pixels.shape[0] != strip.shape[0] \
|
||||
or pixels.shape[2] != strip.shape[2]:
|
||||
return 0
|
||||
width = pixels.shape[1]
|
||||
lo, hi = max(0, int(x)), min(strip.shape[1], int(x) + width)
|
||||
if hi <= lo:
|
||||
return 0
|
||||
strip[:, lo:hi] = pixels[:, lo - int(x):hi - int(x)]
|
||||
if self.__dict__.get('_cached_image') is not None \
|
||||
or self.__dict__.get('_image_source') is not None:
|
||||
self._defer_image()
|
||||
return (hi - lo) * strip.shape[0] * strip.shape[2]
|
||||
|
||||
def remaining_unscrolled(self) -> int:
|
||||
"""Columns of strip still to the right of the viewport."""
|
||||
if not self.has_strip():
|
||||
|
||||
@@ -915,8 +915,7 @@ class DisplayManager:
|
||||
``display.dirty_tracking: false`` if a redraw issue is ever suspected.
|
||||
|
||||
Serialized via ``_update_lock``: plugins can call this directly from
|
||||
background threads (e.g. sports base classes push an immediate
|
||||
"live" refresh from inside update()), so without a lock two callers
|
||||
background threads of their own, so without a lock two callers
|
||||
could both pass the digest check before either writes it back,
|
||||
double-pushing a frame, or interleave the offscreen/current canvas
|
||||
swap below. The lock is scoped to this method, so callers never
|
||||
|
||||
@@ -67,7 +67,7 @@ def _as_vegas_canvas(plugin: Any, display_manager: Any, width: int) -> Iterator[
|
||||
|
||||
|
||||
def render_vegas_elements(plugin: Any, display_manager: Any,
|
||||
width: Optional[int] = None) -> Optional[list]:
|
||||
width: Optional[int] = None) -> Any:
|
||||
"""Call ``plugin.get_vegas_elements()`` as the Vegas ticker does."""
|
||||
render_width = int(width or display_manager.width)
|
||||
with _as_vegas_canvas(plugin, display_manager, render_width):
|
||||
@@ -151,8 +151,9 @@ def check_vegas_elements(plugin: Any, display_manager: Any) -> VegasElementRepor
|
||||
report.errors.append(f"{where} appears twice; keys must be unique")
|
||||
continue
|
||||
seen.add(key)
|
||||
if not isinstance(element.image, Image.Image):
|
||||
report.errors.append(f"{where} image is a {type(element.image).__name__}")
|
||||
image: Any = element.image # typed Image, but a plugin may pass anything
|
||||
if not isinstance(image, Image.Image):
|
||||
report.errors.append(f"{where} image is a {type(image).__name__}")
|
||||
continue
|
||||
if element.image.height != height:
|
||||
report.errors.append(
|
||||
@@ -296,12 +297,12 @@ def check_plugin_vegas_elements(plugin_id: str, plugin_dir: Any, config: dict,
|
||||
if run_update:
|
||||
try:
|
||||
plugin.update()
|
||||
except _TOLERATED_UPDATE_ERRORS as exc:
|
||||
except Exception as exc: # noqa: BLE001 - a plugin's update can raise anything
|
||||
if not isinstance(exc, _TOLERATED_UPDATE_ERRORS):
|
||||
report.errors.append(f"update() raised {exc!r}")
|
||||
return report
|
||||
report.warnings.append(f"update() had no network ({exc!r}); checked "
|
||||
"with whatever data the plugin starts with")
|
||||
except Exception as exc: # noqa: BLE001 - a plugin's update can raise anything
|
||||
report.errors.append(f"update() raised {exc!r}")
|
||||
return report
|
||||
checked = check_vegas_elements(plugin, display_manager)
|
||||
checked.warnings[:0] = report.warnings
|
||||
return checked
|
||||
|
||||
@@ -387,6 +387,9 @@ class VegasModeCoordinator:
|
||||
adapter = getattr(self, 'plugin_adapter', None)
|
||||
if adapter is not None:
|
||||
adapter.live_elements_enabled = active
|
||||
pipeline = getattr(self, 'render_pipeline', None)
|
||||
if pipeline is not None and hasattr(pipeline, 'set_live'):
|
||||
pipeline.set_live(active)
|
||||
plugin_manager = getattr(self, 'plugin_manager', None)
|
||||
add = getattr(plugin_manager, 'add_update_listener', None)
|
||||
remove = getattr(plugin_manager, 'remove_update_listener', None)
|
||||
@@ -405,9 +408,11 @@ class VegasModeCoordinator:
|
||||
"""Update listener: a plugin's data may have changed.
|
||||
|
||||
Runs on the update worker with the plugin's lock held, so it only
|
||||
moves the plugin's epoch on; whatever redraws happen later read it.
|
||||
moves the plugin's epoch on and wakes the live-element worker, which
|
||||
redraws once the lock is free.
|
||||
"""
|
||||
self.live_epochs.bump(plugin_id)
|
||||
self.render_pipeline.notify_live_data(plugin_id)
|
||||
|
||||
def _install_render_gate(self) -> None:
|
||||
"""Gate the prefetch thread on the render thread's swaps; see VegasModeConfig."""
|
||||
@@ -504,6 +509,12 @@ class VegasModeCoordinator:
|
||||
# game still shown as live the next morning.
|
||||
self.render_pipeline.refresh_updated_plugins()
|
||||
|
||||
# Copy any live-element redraws the worker has finished into the
|
||||
# strip, between this frame and the last. A deque check when there
|
||||
# are none.
|
||||
if self.live_active:
|
||||
self.render_pipeline.apply_live_patches()
|
||||
|
||||
# Extend the strip before the scroll can reach its end, so the next
|
||||
# group arrives from the right and motion never stops. No cycle
|
||||
# boundary, so no freeze, no substitution and no restart with the
|
||||
@@ -715,7 +726,12 @@ class VegasModeCoordinator:
|
||||
# main loop's _tick_plugin_updates() finds all intervals already
|
||||
# satisfied on return, so the inter-iteration gap is <1 ms and the
|
||||
# display never shows a frozen frame between iterations.
|
||||
_UPDATE_TICK_FRAMES = max(1, int(self.render_pipeline.target_fps * 4)) # every 4 s regardless of FPS
|
||||
# Every 4 s, or every 1 s while the strip holds live elements:
|
||||
# plugins are only scheduled on this tick, so its period is added
|
||||
# to how late a live update can be.
|
||||
tick_seconds = (1.0 if self.live_active
|
||||
and self.render_pipeline.has_live_records() else 4.0)
|
||||
_UPDATE_TICK_FRAMES = max(1, int(self.render_pipeline.target_fps * tick_seconds))
|
||||
if (self._update_callback and
|
||||
frame_count % _UPDATE_TICK_FRAMES == 0 and
|
||||
not self._update_tick_running):
|
||||
|
||||
@@ -60,6 +60,41 @@ class ElementRecord(NamedTuple):
|
||||
refresh_hz: float
|
||||
|
||||
|
||||
class RenderedElement(NamedTuple):
|
||||
"""One live element freshly redrawn by the worker, ready to compare and swap."""
|
||||
key: str
|
||||
#: The plugin's data epoch it was drawn from.
|
||||
epoch: int
|
||||
version: object
|
||||
#: Pinned pixels (see pin_element), read-only.
|
||||
pixels: np.ndarray
|
||||
digest: Tuple[Tuple[int, ...], int]
|
||||
#: Pinned width, the width it would occupy in the strip.
|
||||
width: int
|
||||
|
||||
|
||||
class LivePatch(NamedTuple):
|
||||
"""A redraw handed from the worker to the render thread for one record."""
|
||||
seq: int
|
||||
#: The strip generation it was made against; a patch for an older strip
|
||||
#: is dropped.
|
||||
strip_gen: int
|
||||
epoch: int
|
||||
pixels: np.ndarray
|
||||
digest: Tuple[Tuple[int, ...], int]
|
||||
made_at: float
|
||||
|
||||
|
||||
class LiveView(NamedTuple):
|
||||
"""Where the viewport is, in absolute strip columns, published every frame."""
|
||||
abs_left: int
|
||||
abs_right: int
|
||||
#: The end of the strip: how far ahead content exists.
|
||||
abs_end: int
|
||||
#: time.monotonic() when published. An old one means frames have stopped.
|
||||
t_mono: float
|
||||
|
||||
|
||||
def tag(image: Image.Image, meta: ElementMeta) -> Image.Image:
|
||||
"""Mark ``image`` as the live element ``meta`` describes. Returns it."""
|
||||
image.info[INFO_KEY] = meta
|
||||
|
||||
@@ -0,0 +1,503 @@
|
||||
"""The one background worker behind live Vegas elements.
|
||||
|
||||
Vegas draws everything the strip shows off the render thread. Until live
|
||||
elements that was one short-lived prefetch thread per group; now, once the
|
||||
strip holds a live element, it is this worker, which does three kinds of job
|
||||
one at a time, most urgent first:
|
||||
|
||||
- **group steps**: fetching the next group of plugins for the strip, one
|
||||
plugin per step (what the prefetch thread did in one go);
|
||||
- **data refreshes**: when a plugin's data has moved on (its epoch, see
|
||||
elements.LiveEpochs) past what its elements in the strip were drawn from,
|
||||
redraw them and hand over the ones whose pixels changed;
|
||||
- **ticks**: redraw an element that animates (``refresh_hz``) while it is on
|
||||
or near the screen.
|
||||
|
||||
Nothing here touches the strip. A finished redraw becomes a
|
||||
:class:`~src.vegas_mode.elements.LivePatch` in the pipeline's slot for that
|
||||
element (one per element, the latest wins) and the render thread copies it
|
||||
into the strip between two frames (RenderPipeline.apply_live_patches). The
|
||||
hand-over is lock-free: a dict store and a deque append here, a deque popleft
|
||||
and a dict pop there, so the render thread never waits on this thread.
|
||||
|
||||
Every job runs inside the render gate when there is one (src/common/
|
||||
render_gate.py), so Python runs here only while the render thread is waiting
|
||||
for the panel. Without it (the stock rgbmatrix binding) animation is capped at
|
||||
:data:`UNGATED_MAX_HZ`.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import collections
|
||||
import logging
|
||||
import os
|
||||
import queue
|
||||
import threading
|
||||
import time
|
||||
from contextlib import nullcontext
|
||||
from typing import Any, Callable, Dict, List, Optional, Set, Tuple
|
||||
|
||||
from src.vegas_mode.elements import ElementRecord, LivePatch, LiveView
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
#: A view older than this means frames have stopped (a paused scroll, an
|
||||
#: interrupt): only group work runs, since nothing redrawn would be seen.
|
||||
VIEW_STALE_S = 0.5
|
||||
#: How long a data refresh waits for the plugin's lock before trying later.
|
||||
DATA_LOCK_TIMEOUT = 0.25
|
||||
#: ...and how much later.
|
||||
LOCK_BACKOFF_S = 1.0
|
||||
#: Animation ceiling without the render gate, where every redraw competes
|
||||
#: with the render thread for the GIL.
|
||||
UNGATED_MAX_HZ = 1.0
|
||||
#: An element whose redraws take longer than this on average is animated at
|
||||
#: half its rate, down to MIN_THROTTLED_HZ.
|
||||
SLOW_RENDER_S = 0.05
|
||||
MIN_THROTTLED_HZ = 0.5
|
||||
#: Longest the worker sleeps with nothing due, so a floor or a backoff that
|
||||
#: expires is noticed.
|
||||
IDLE_WAIT_S = 0.5
|
||||
#: How often the worker logs what it did, when it did anything.
|
||||
SUMMARY_INTERVAL_S = 300.0
|
||||
#: Weight of the newest sample in the per-element render time average.
|
||||
EWMA_ALPHA = 0.2
|
||||
#: How long the worker waits for a one-shot prefetch thread it takes over from.
|
||||
LEGACY_PREFETCH_JOIN_S = 15.0
|
||||
|
||||
|
||||
def _visible(record: ElementRecord, view: LiveView) -> bool:
|
||||
return record.abs_x < view.abs_right and record.abs_x + record.width > view.abs_left
|
||||
|
||||
|
||||
def _behind(record: ElementRecord, view: LiveView) -> bool:
|
||||
return record.abs_x + record.width <= view.abs_left
|
||||
|
||||
|
||||
class _GroupJob:
|
||||
"""A group fetch in progress, one member per step."""
|
||||
|
||||
def __init__(self, generation: int, plugin_ids: List[str]) -> None:
|
||||
self.generation = generation
|
||||
self.pending = list(plugin_ids)
|
||||
self.group: List[Tuple[str, Any]] = []
|
||||
|
||||
|
||||
class VegasWorker(threading.Thread):
|
||||
"""See the module docstring. Owned by the RenderPipeline that starts it."""
|
||||
|
||||
def __init__(self, pipeline: Any, clock: Callable[[], float] = time.monotonic) -> None:
|
||||
super().__init__(daemon=True, name="vegas-live-worker")
|
||||
self.pipeline = pipeline
|
||||
#: time.monotonic, or a fake one in tests; the pipeline's view is
|
||||
#: stamped with time.monotonic too.
|
||||
self._clock = clock
|
||||
self.inbox: "queue.SimpleQueue[Tuple[str, Any]]" = queue.SimpleQueue()
|
||||
self._stopping = False
|
||||
self._group_wanted = False
|
||||
self._group_job: Optional[_GroupJob] = None
|
||||
self._strip_gen = pipeline._strip_gen
|
||||
# Per element (record seq): the epoch this worker last handed over,
|
||||
# or found needed nothing; the digest of its latest hand-over; when
|
||||
# its next animation tick is due.
|
||||
self._handled_epoch: Dict[int, int] = {}
|
||||
self._handed_digest: Dict[int, Any] = {}
|
||||
self._next_tick: Dict[int, float] = {}
|
||||
# Per plugin: when its last data refresh ran, and a lock backoff.
|
||||
self._last_data_job: Dict[str, float] = {}
|
||||
self._backoff_until: Dict[str, float] = {}
|
||||
# Per (plugin, key): average redraw time, for throttling.
|
||||
self._render_ewma: Dict[Tuple[str, str], float] = {}
|
||||
self._refused: Set[Tuple[str, str, int]] = set()
|
||||
self.stats: collections.Counter = collections.Counter()
|
||||
self.busy_seconds = 0.0
|
||||
self._began = clock()
|
||||
self._last_summary = self._began
|
||||
self._jobs_since_prune = 0
|
||||
# The slowest redraw since the last summary: (seconds, (plugin, key)).
|
||||
self._slowest_redraw: Optional[Tuple[float, Tuple[str, str]]] = None
|
||||
|
||||
# -- control, from other threads ------------------------------------------
|
||||
|
||||
def request_group(self) -> None:
|
||||
"""Fetch the next group for the strip when nothing more urgent is due."""
|
||||
self._group_wanted = True
|
||||
self.inbox.put(("group", None))
|
||||
|
||||
def notify_data(self, plugin_id: str) -> None:
|
||||
"""A plugin's data moved on. Only a wake-up: its epoch is the truth."""
|
||||
self.inbox.put(("data", plugin_id))
|
||||
|
||||
def stop(self) -> None:
|
||||
"""Stop after the current job. Does not wait for it."""
|
||||
self._stopping = True
|
||||
self.inbox.put(("stop", None))
|
||||
|
||||
# -- the loop -------------------------------------------------------------
|
||||
|
||||
def run(self) -> None:
|
||||
try:
|
||||
# Linux applies nice per thread: deprioritise against the render
|
||||
# loop, as the one-shot prefetch thread always did.
|
||||
os.nice(10)
|
||||
except (OSError, AttributeError):
|
||||
pass
|
||||
self._join_legacy_prefetch()
|
||||
while not self._should_stop():
|
||||
self._wait(self._next_wait(self._clock()))
|
||||
if self._should_stop():
|
||||
break
|
||||
job = self._pick(self._clock())
|
||||
if job is not None:
|
||||
self._run(job)
|
||||
self._maybe_summarise()
|
||||
self._hand_over_partial_group()
|
||||
logger.debug("Vegas live worker stopped")
|
||||
|
||||
def _should_stop(self) -> bool:
|
||||
# A method, not a bare attribute read: stop() sets it from another
|
||||
# thread between two reads in run().
|
||||
return self._stopping
|
||||
|
||||
def _join_legacy_prefetch(self) -> None:
|
||||
"""Let a one-shot prefetch thread, or a worker stopped earlier, finish first.
|
||||
|
||||
So that only one thread ever draws for the strip: measured on hdpi,
|
||||
each extra thread competing for the GIL made the render thread late
|
||||
more often, not less (src/common/render_gate.py).
|
||||
"""
|
||||
for name in ('_prefetch_thread', '_retired_worker'):
|
||||
thread = getattr(self.pipeline, name, None)
|
||||
if thread is not None and thread is not self \
|
||||
and thread is not threading.current_thread() and thread.is_alive():
|
||||
thread.join(LEGACY_PREFETCH_JOIN_S)
|
||||
|
||||
def _wait(self, timeout: float) -> None:
|
||||
try:
|
||||
message = self.inbox.get(timeout=max(0.0, timeout))
|
||||
except queue.Empty:
|
||||
return
|
||||
while True:
|
||||
if message[0] == "stop":
|
||||
self._stopping = True
|
||||
elif message[0] == "group":
|
||||
self._group_wanted = True
|
||||
try:
|
||||
message = self.inbox.get_nowait()
|
||||
except queue.Empty:
|
||||
return
|
||||
|
||||
def _next_wait(self, now: float) -> float:
|
||||
"""Seconds until something may be due: the next tick, else IDLE_WAIT_S."""
|
||||
p = self.pipeline
|
||||
if self._group_job is not None or (
|
||||
self._group_wanted and p._prepared_group is None):
|
||||
return 0.0
|
||||
# Ticks only count while frames are flowing: during a pause _pick runs
|
||||
# none, and a past-due tick would otherwise make this 0 and spin.
|
||||
view = p._view
|
||||
if view is None or now - view.t_mono > VIEW_STALE_S:
|
||||
return IDLE_WAIT_S
|
||||
# Only elements _due_tick would run. A tick left behind by an element
|
||||
# trimmed away, or one no longer animated, is never run, and counting
|
||||
# it held this at its floor: a spin at 100 wake-ups a second.
|
||||
soonest = now + IDLE_WAIT_S
|
||||
lead = self._tick_lead()
|
||||
for record in p._elements:
|
||||
if self._tickable(record, view, lead):
|
||||
soonest = min(soonest, self._next_tick.get(record.seq, now))
|
||||
return max(0.01, soonest - now)
|
||||
|
||||
# -- choosing ---------------------------------------------------------------
|
||||
|
||||
def _pick(self, now: float) -> Optional[Tuple[str, Any]]:
|
||||
"""The most urgent job, or None. See the module docstring for the order."""
|
||||
p = self.pipeline
|
||||
if p._strip_gen != self._strip_gen:
|
||||
self._forget_everything(p._strip_gen)
|
||||
view = p._view
|
||||
fresh = view is not None and now - view.t_mono <= VIEW_STALE_S
|
||||
group_ready = self._group_job is not None or (
|
||||
self._group_wanted and p._prepared_group is None)
|
||||
records = p._elements
|
||||
|
||||
if group_ready and fresh and self._group_urgent(view):
|
||||
return ("group", None)
|
||||
if fresh and records:
|
||||
plugin_id = self._due_data(records, view, now, visible_only=True)
|
||||
if plugin_id is not None:
|
||||
return ("data", plugin_id)
|
||||
record = self._due_tick(records, view, now)
|
||||
if record is not None:
|
||||
return ("tick", record)
|
||||
if group_ready:
|
||||
return ("group", None)
|
||||
if fresh and records:
|
||||
plugin_id = self._due_data(records, view, now, visible_only=False)
|
||||
if plugin_id is not None:
|
||||
return ("data", plugin_id)
|
||||
return None
|
||||
|
||||
def _group_urgent(self, view: LiveView) -> bool:
|
||||
width = self.pipeline.display_width
|
||||
threshold = (self.pipeline.config.extend_threshold_screens + 1.0) * width
|
||||
return bool(view.abs_end - view.abs_right <= threshold)
|
||||
|
||||
def _epoch(self, plugin_id: str) -> int:
|
||||
epochs = getattr(self.pipeline.stream_manager.plugin_adapter, 'live_epochs', None)
|
||||
return int(epochs.get(plugin_id)) if epochs is not None else 0
|
||||
|
||||
def _done_epoch(self, record: ElementRecord) -> int:
|
||||
applied = self.pipeline._applied.get(record.seq)
|
||||
return max(applied[0] if applied is not None else record.epoch,
|
||||
self._handled_epoch.get(record.seq, -1))
|
||||
|
||||
def _due_data(self, records: Tuple[ElementRecord, ...], view: LiveView,
|
||||
now: float, visible_only: bool) -> Optional[str]:
|
||||
"""The plugin whose stale elements are nearest the screen, if any may redraw."""
|
||||
floor = self.pipeline.config.live_min_interval
|
||||
best: Optional[Tuple[int, str]] = None
|
||||
for record in records:
|
||||
if _behind(record, view):
|
||||
continue
|
||||
if visible_only and not _visible(record, view):
|
||||
continue
|
||||
plugin_id = record.plugin_id
|
||||
if self._epoch(plugin_id) <= self._done_epoch(record):
|
||||
continue
|
||||
if now < self._backoff_until.get(plugin_id, 0.0):
|
||||
continue
|
||||
if now - self._last_data_job.get(plugin_id, float('-inf')) < floor:
|
||||
continue
|
||||
distance = max(0, record.abs_x - view.abs_right)
|
||||
if best is None or distance < best[0]:
|
||||
best = (distance, plugin_id)
|
||||
return best[1] if best is not None else None
|
||||
|
||||
def _tick_hz(self, record: ElementRecord) -> float:
|
||||
cfg = self.pipeline.config
|
||||
hz = min(float(record.refresh_hz), float(cfg.live_max_hz))
|
||||
if getattr(self.pipeline.display_manager, 'render_gate', None) is None:
|
||||
hz = min(hz, UNGATED_MAX_HZ)
|
||||
ewma = self._render_ewma.get((record.plugin_id, record.key), 0.0)
|
||||
if hz > 0 and ewma > SLOW_RENDER_S:
|
||||
# Halved, but never below the floor -- nor raised to it, for an
|
||||
# element already asking for less.
|
||||
hz = min(hz, max(MIN_THROTTLED_HZ, hz / 2.0))
|
||||
return hz
|
||||
|
||||
def _tick_lead(self) -> float:
|
||||
return float(self.pipeline.config.live_lead_screens * self.pipeline.display_width)
|
||||
|
||||
def _tickable(self, record: ElementRecord, view: LiveView, lead: float) -> bool:
|
||||
"""Whether an element animates now: it has a rate, and is on or near the screen."""
|
||||
return (record.refresh_hz > 0 and self._tick_hz(record) > 0
|
||||
and not _behind(record, view) and record.abs_x < view.abs_right + lead)
|
||||
|
||||
def _due_tick(self, records: Tuple[ElementRecord, ...], view: LiveView,
|
||||
now: float) -> Optional[ElementRecord]:
|
||||
lead = self._tick_lead()
|
||||
best: Optional[ElementRecord] = None
|
||||
best_due = 0.0
|
||||
for record in records:
|
||||
if not self._tickable(record, view, lead):
|
||||
self._next_tick.pop(record.seq, None)
|
||||
continue
|
||||
due = self._next_tick.get(record.seq, now)
|
||||
if due <= now and (best is None or due < best_due):
|
||||
best, best_due = record, due
|
||||
return best
|
||||
|
||||
# -- running ----------------------------------------------------------------
|
||||
|
||||
def _run(self, job: Tuple[str, Any]) -> None:
|
||||
kind, arg = job
|
||||
gate = getattr(self.pipeline.display_manager, 'render_gate', None)
|
||||
started = self._clock()
|
||||
try:
|
||||
with gate.yielding() if gate is not None else nullcontext():
|
||||
if kind == "group":
|
||||
self._group_step()
|
||||
elif kind == "data":
|
||||
self._data_job(arg, started)
|
||||
else:
|
||||
self._tick_job(arg, started)
|
||||
self.stats[kind] += 1
|
||||
except Exception as exc: # pylint: disable=broad-except
|
||||
# Plugin code runs in here; one bad job must not end the worker,
|
||||
# which also fetches every group for the strip.
|
||||
self.stats["errors"] += 1
|
||||
if self.stats["errors"] <= 3 or self.stats["errors"] % 100 == 0:
|
||||
logger.exception("Vegas live worker: %s job failed (%s)", kind, exc)
|
||||
finally:
|
||||
self.busy_seconds += self._clock() - started
|
||||
self._jobs_since_prune += 1
|
||||
if self._jobs_since_prune >= 64:
|
||||
self._prune()
|
||||
|
||||
def _hand_over_partial_group(self) -> None:
|
||||
"""On stopping: publish the members of a group already fetched.
|
||||
|
||||
Its plugins were taken from the rotation when it was planned, so a
|
||||
group dropped here would skip them until the next cycle. The rest of
|
||||
it is not fetched. Nothing is published over a group already waiting,
|
||||
or into a Vegas reset since.
|
||||
"""
|
||||
job, self._group_job = self._group_job, None
|
||||
if job is None or not job.group:
|
||||
return
|
||||
p = self.pipeline
|
||||
with p._prefetch_lock:
|
||||
if job.generation == p._prefetch_generation and p._prepared_group is None:
|
||||
p._prepared_group = job.group
|
||||
|
||||
def _group_step(self) -> None:
|
||||
p = self.pipeline
|
||||
job = self._group_job
|
||||
if job is None:
|
||||
with p._prefetch_lock:
|
||||
if p._prepared_group is not None:
|
||||
self._group_wanted = False
|
||||
return
|
||||
generation = p._prefetch_generation
|
||||
self._group_wanted = False
|
||||
job = self._group_job = _GroupJob(generation, p.stream_manager.plan_next_group())
|
||||
if job.generation != p._prefetch_generation:
|
||||
self._group_job = None # Vegas was reset meanwhile
|
||||
return
|
||||
if job.pending:
|
||||
member = p.stream_manager.fetch_group_member(
|
||||
job.pending.pop(0), offscreen_only=True)
|
||||
if member is not None:
|
||||
job.group.append(member)
|
||||
if not job.pending:
|
||||
self._group_job = None
|
||||
with p._prefetch_lock:
|
||||
if job.generation == p._prefetch_generation:
|
||||
p._prepared_group = job.group
|
||||
|
||||
def _data_job(self, plugin_id: str, now: float) -> None:
|
||||
p = self.pipeline
|
||||
self._last_data_job[plugin_id] = now
|
||||
plugin = getattr(p.stream_manager.plugin_manager, 'plugins', {}).get(plugin_id)
|
||||
if plugin is None:
|
||||
return
|
||||
gen = p._strip_gen
|
||||
batch = p.stream_manager.plugin_adapter.render_live_elements(
|
||||
plugin, plugin_id, lock_timeout=DATA_LOCK_TIMEOUT)
|
||||
if batch is None:
|
||||
self.stats["lock_busy"] += 1
|
||||
self._backoff_until[plugin_id] = now + LOCK_BACKOFF_S
|
||||
return
|
||||
epoch, rendered = batch
|
||||
view = p._view
|
||||
for record in p._elements:
|
||||
if record.plugin_id != plugin_id:
|
||||
continue
|
||||
if view is not None and _behind(record, view):
|
||||
continue
|
||||
element = rendered.get(record.key)
|
||||
if element is not None:
|
||||
self._hand_over(record, element, epoch, gen)
|
||||
# A key the plugin no longer has keeps its last pixels until it
|
||||
# scrolls off; either way this epoch is dealt with.
|
||||
self._handled_epoch[record.seq] = max(
|
||||
epoch, self._handled_epoch.get(record.seq, -1))
|
||||
|
||||
def _tick_job(self, record: ElementRecord, now: float) -> None:
|
||||
p = self.pipeline
|
||||
hz = self._tick_hz(record)
|
||||
self._next_tick[record.seq] = now + (1.0 / hz if hz > 0 else IDLE_WAIT_S)
|
||||
plugin = getattr(p.stream_manager.plugin_manager, 'plugins', {}).get(record.plugin_id)
|
||||
if plugin is None:
|
||||
return
|
||||
adapter = p.stream_manager.plugin_adapter
|
||||
key = (record.plugin_id, record.key)
|
||||
at = now + self._render_ewma.get(key, 0.0) + p.frame_interval
|
||||
started = time.perf_counter()
|
||||
if adapter.has_lock_free_redraw(plugin):
|
||||
# None from the plugin means nothing to redraw this time.
|
||||
element = adapter.redraw_live_element(
|
||||
plugin, record.plugin_id, record.key, record.width, p.display_height, at)
|
||||
else:
|
||||
# No lock-free redraw: redraw everything, but never wait for the
|
||||
# plugin's lock (update() may be doing network I/O under it).
|
||||
batch = adapter.render_live_elements(plugin, record.plugin_id, lock_timeout=0.0)
|
||||
element = batch[1].get(record.key) if batch is not None else None
|
||||
took = time.perf_counter() - started
|
||||
if self._slowest_redraw is None or took > self._slowest_redraw[0]:
|
||||
self._slowest_redraw = (took, key)
|
||||
previous = self._render_ewma.get(key)
|
||||
self._render_ewma[key] = took if previous is None else (
|
||||
EWMA_ALPHA * took + (1.0 - EWMA_ALPHA) * previous)
|
||||
if element is not None:
|
||||
self._hand_over(record, element, element.epoch, p._strip_gen)
|
||||
|
||||
def _hand_over(self, record: ElementRecord, element: Any, epoch: int, gen: int) -> None:
|
||||
"""Queue a redraw for the render thread, unless nothing would change."""
|
||||
p = self.pipeline
|
||||
if element.width != record.width:
|
||||
marker = (record.plugin_id, record.key, element.width)
|
||||
if marker not in self._refused:
|
||||
self._refused.add(marker)
|
||||
logger.info(
|
||||
"[%s] Live element %r redrawn %dpx wide, placed at %dpx; "
|
||||
"kept as it was (a live element's width must not change)",
|
||||
record.plugin_id, record.key, element.width, record.width)
|
||||
self.stats["refused"] += 1
|
||||
return
|
||||
last = self._handed_digest.get(record.seq)
|
||||
if last is None:
|
||||
applied = p._applied.get(record.seq)
|
||||
last = applied[1] if applied is not None else record.digest
|
||||
if element.digest == last:
|
||||
self.stats["unchanged"] += 1
|
||||
return
|
||||
p._live_slots[record.seq] = LivePatch(
|
||||
seq=record.seq, strip_gen=gen, epoch=epoch, pixels=element.pixels,
|
||||
digest=element.digest, made_at=self._clock())
|
||||
p._live_ready.append(record.seq)
|
||||
self._handed_digest[record.seq] = element.digest
|
||||
self.stats["patches"] += 1
|
||||
|
||||
# -- housekeeping -----------------------------------------------------------
|
||||
|
||||
def _forget_everything(self, gen: int) -> None:
|
||||
self._strip_gen = gen
|
||||
self._handled_epoch.clear()
|
||||
self._handed_digest.clear()
|
||||
self._next_tick.clear()
|
||||
|
||||
def _prune(self) -> None:
|
||||
self._jobs_since_prune = 0
|
||||
live = {record.seq for record in self.pipeline._elements}
|
||||
for table in (self._handled_epoch, self._handed_digest, self._next_tick):
|
||||
for seq in [s for s in table if s not in live]:
|
||||
del table[seq]
|
||||
|
||||
def _maybe_summarise(self) -> None:
|
||||
now = self._clock()
|
||||
if now - self._last_summary < SUMMARY_INTERVAL_S:
|
||||
return
|
||||
elapsed = now - self._last_summary
|
||||
self._last_summary = now
|
||||
stats, self.stats = self.stats, collections.Counter()
|
||||
busy, self.busy_seconds = self.busy_seconds, 0.0
|
||||
slowest, self._slowest_redraw = self._slowest_redraw, None
|
||||
if not stats:
|
||||
return
|
||||
redraws = ""
|
||||
if slowest is not None:
|
||||
# What a tick costs is the number that decides whether an
|
||||
# animated element can keep its rate: say it for the worst one.
|
||||
took, (plugin_id, key) = slowest
|
||||
average = self._render_ewma.get((plugin_id, key), took)
|
||||
redraws = "; slowest redraw %.1fms (%s %r, average %.1fms)" % (
|
||||
took * 1000.0, plugin_id, key, average * 1000.0)
|
||||
logger.info(
|
||||
"Vegas live: %d group step(s), %d data refresh(es), %d tick(s); "
|
||||
"%d patch(es) handed over, %d unchanged, %d refused, %d lock-busy, "
|
||||
"%d error(s); worker busy %.1f%%%s",
|
||||
stats["group"], stats["data"], stats["tick"], stats["patches"],
|
||||
stats["unchanged"], stats["refused"], stats["lock_busy"],
|
||||
stats["errors"], 100.0 * busy / elapsed if elapsed else 0.0, redraws)
|
||||
@@ -9,7 +9,7 @@ import logging
|
||||
import threading
|
||||
import time
|
||||
from contextlib import contextmanager, nullcontext
|
||||
from typing import Optional, List, Any, Tuple, Union, TYPE_CHECKING
|
||||
from typing import Dict, Optional, List, Any, Tuple, Union, TYPE_CHECKING
|
||||
from PIL import Image
|
||||
|
||||
from src.common.scroll_helper import ScrollHelper
|
||||
@@ -18,6 +18,7 @@ from src.plugin_system.vegas_elements import VegasElement
|
||||
from src.vegas_mode.elements import (
|
||||
ElementMeta,
|
||||
LiveEpochs,
|
||||
RenderedElement,
|
||||
meta_of,
|
||||
pin_element,
|
||||
pixel_digest,
|
||||
@@ -109,6 +110,11 @@ class PluginAdapter:
|
||||
# Element problems already reported, so a plugin with a bad hook logs
|
||||
# once rather than on every fetch.
|
||||
self._element_warnings: set = set()
|
||||
# (plugin_id, key) -> (version, source image, (padding, height),
|
||||
# pinned pixels, digest) of the last conversion, so an element handed
|
||||
# back unchanged is not converted again. Only the live-element worker
|
||||
# reads or writes it; invalidate_cache() swaps in a fresh one.
|
||||
self._element_memo: Dict[Tuple[str, str], Tuple[Any, ...]] = {}
|
||||
|
||||
logger.debug(
|
||||
"PluginAdapter initialized: display=%dx%d",
|
||||
@@ -198,18 +204,20 @@ class PluginAdapter:
|
||||
return bool(raw)
|
||||
|
||||
@contextmanager
|
||||
def _plugin_lock(self, plugin_id: str):
|
||||
def _plugin_lock(self, plugin_id: str, timeout: Optional[float] = None):
|
||||
"""Hold the plugin's update/display lock, waiting a bounded time.
|
||||
|
||||
Yields whether it was acquired. Yields True, holding nothing, when
|
||||
there is no plugin manager to ask -- the behaviour before the lock was
|
||||
taken here at all.
|
||||
taken here at all. ``timeout`` defaults to PLUGIN_LOCK_TIMEOUT; 0
|
||||
does not wait at all.
|
||||
"""
|
||||
if not hasattr(self.plugin_manager, 'get_plugin_lock'):
|
||||
yield True
|
||||
return
|
||||
lock = self.plugin_manager.get_plugin_lock(plugin_id)
|
||||
acquired = lock.acquire(timeout=self.PLUGIN_LOCK_TIMEOUT)
|
||||
wait = self.PLUGIN_LOCK_TIMEOUT if timeout is None else timeout
|
||||
acquired = lock.acquire(timeout=wait) if wait > 0 else lock.acquire(blocking=False)
|
||||
try:
|
||||
yield acquired
|
||||
finally:
|
||||
@@ -455,18 +463,21 @@ class PluginAdapter:
|
||||
if isinstance(plugin_cfg, dict):
|
||||
raw = plugin_cfg.get('vegas_width_pct')
|
||||
if raw not in (None, ''):
|
||||
# Reported once per value: the live paths resolve the width on
|
||||
# every redraw, several times a second for an animated element.
|
||||
try:
|
||||
candidate = int(raw)
|
||||
except (TypeError, ValueError):
|
||||
logger.warning(
|
||||
"[%s] Invalid vegas_width_pct %r, ignoring", plugin_id, raw)
|
||||
self._warn_element_once(
|
||||
plugin_id, "Invalid vegas_width_pct %r, ignoring", raw,
|
||||
once_key=repr(raw))
|
||||
else:
|
||||
if 10 <= candidate <= 100:
|
||||
pct = candidate
|
||||
else:
|
||||
logger.warning(
|
||||
"[%s] vegas_width_pct %d out of range 10-100, ignoring",
|
||||
plugin_id, candidate)
|
||||
self._warn_element_once(
|
||||
plugin_id, "vegas_width_pct %d out of range 10-100, ignoring",
|
||||
candidate, once_key=repr(raw))
|
||||
|
||||
if pct >= 100:
|
||||
return self.display_width
|
||||
@@ -675,8 +686,7 @@ class PluginAdapter:
|
||||
|
||||
if len(images) == 1:
|
||||
only = images[0]
|
||||
pad = self.config.content_padding if self.config.auto_trim else 0
|
||||
if meta_of(only) is not None and only.width - 2 * pad <= budget:
|
||||
if meta_of(only) is not None and only.width - 2 * self._padding() <= budget:
|
||||
# A live element's pinned margins are not content. One whose
|
||||
# drawing fits the budget is kept whole, and live, rather than
|
||||
# cut for the sake of its own blank padding.
|
||||
@@ -840,9 +850,14 @@ class PluginAdapter:
|
||||
)
|
||||
return img.crop((start, 0, end, img.height))
|
||||
|
||||
def _warn_element_once(self, plugin_id: str, problem: str, *args: Any) -> None:
|
||||
"""Report a plugin's element problem once per process, then quietly."""
|
||||
key = (plugin_id, problem)
|
||||
def _warn_element_once(self, plugin_id: str, problem: str, *args: Any,
|
||||
once_key: Optional[str] = None) -> None:
|
||||
"""Report a plugin's problem once per process, then quietly.
|
||||
|
||||
Once per ``problem`` (the format string), or per ``once_key`` within it
|
||||
when given, so a different bad value is still reported.
|
||||
"""
|
||||
key = (plugin_id, problem, once_key)
|
||||
if key in self._element_warnings:
|
||||
logger.debug("[%s] " + problem, plugin_id, *args)
|
||||
return
|
||||
@@ -885,16 +900,98 @@ class PluginAdapter:
|
||||
"be used (%r); using its get_vegas_content() instead", exc)
|
||||
return None
|
||||
|
||||
def _images_from_elements(
|
||||
self, result: Any, plugin_id: str, epoch: int
|
||||
) -> Optional[List[Image.Image]]:
|
||||
"""Turn get_vegas_elements()'s answer into images for the pipeline.
|
||||
def render_live_elements(
|
||||
self, plugin: 'BasePlugin', plugin_id: str, lock_timeout: float
|
||||
) -> Optional[Tuple[int, Dict[str, RenderedElement]]]:
|
||||
"""Redraw a plugin's live elements for the live-element worker.
|
||||
|
||||
Live elements come out pinned (RGB, display height, content_padding
|
||||
black each side, never trimmed afterwards) and tagged with their
|
||||
ElementMeta; plain ones (``live=False``) come out as ordinary content.
|
||||
Anything that is not a usable element is dropped with a warning; a
|
||||
duplicate key keeps its first element.
|
||||
Like the keyed fetch, but for a strip that already holds the elements:
|
||||
no cache (the caller knows the plugin's data moved on), and each live
|
||||
element comes back as a RenderedElement to compare with what the strip
|
||||
shows. An element whose ``version`` is the one already redrawn reuses
|
||||
its pinned pixels and digest, so an unchanged scoreboard costs the
|
||||
plugin's own version check and no conversion.
|
||||
|
||||
Returns ``(epoch, {key: element})``, the epoch read under the lock; an
|
||||
empty dict when the plugin had nothing (or failed, logged once). None
|
||||
only when the lock could not be had within ``lock_timeout`` -- the
|
||||
caller tries again later.
|
||||
"""
|
||||
with self._plugin_lock(plugin_id, timeout=lock_timeout) as acquired:
|
||||
if not acquired:
|
||||
return None
|
||||
epochs = self.live_epochs
|
||||
epoch = epochs.get(plugin_id) if epochs is not None else 0
|
||||
render_width = self.resolve_render_width(plugin, plugin_id)
|
||||
plugin._vegas_render_width = render_width
|
||||
try:
|
||||
with self._isolated_canvas(render_width):
|
||||
result = plugin.get_vegas_elements()
|
||||
except Exception as exc: # pylint: disable=broad-except
|
||||
self._warn_element_once(
|
||||
plugin_id, "get_vegas_elements() raised %r while redrawing; "
|
||||
"its elements keep what they show", exc)
|
||||
return epoch, {}
|
||||
finally:
|
||||
plugin._vegas_render_width = None
|
||||
return epoch, self._rendered_from_elements(result, plugin_id, epoch)
|
||||
|
||||
@staticmethod
|
||||
def has_lock_free_redraw(plugin: Any) -> bool:
|
||||
"""Whether the plugin's class overrides BasePlugin.redraw_vegas_element."""
|
||||
method = getattr(type(plugin), 'redraw_vegas_element', None)
|
||||
return method is not None and method is not _BasePlugin.redraw_vegas_element
|
||||
|
||||
def redraw_live_element(
|
||||
self, plugin: 'BasePlugin', plugin_id: str, key: str, width: int,
|
||||
height: int, at: float
|
||||
) -> Optional[RenderedElement]:
|
||||
"""One element redrawn for a moment in time, without the plugin's lock.
|
||||
|
||||
``width`` is the element's width in the strip (pinned); the plugin is
|
||||
asked for that less its padding, exactly, and anything else is
|
||||
refused. None when the plugin has no lock-free redraw, returns None,
|
||||
or fails (logged once).
|
||||
"""
|
||||
if not self.has_lock_free_redraw(plugin):
|
||||
return None
|
||||
padding = self._padding()
|
||||
inner = width - 2 * padding
|
||||
if inner <= 0:
|
||||
return None
|
||||
epochs = self.live_epochs
|
||||
epoch = epochs.get(plugin_id) if epochs is not None else 0
|
||||
render_width = self.resolve_render_width(plugin, plugin_id)
|
||||
plugin._vegas_render_width = render_width
|
||||
try:
|
||||
with self._isolated_canvas(render_width):
|
||||
image = plugin.redraw_vegas_element(key, inner, height, at)
|
||||
except Exception as exc: # pylint: disable=broad-except
|
||||
self._warn_element_once(
|
||||
plugin_id, "redraw_vegas_element(%r) raised %r", key, exc)
|
||||
return None
|
||||
finally:
|
||||
plugin._vegas_render_width = None
|
||||
if image is None:
|
||||
return None
|
||||
if not isinstance(image, Image.Image) or image.size != (inner, height):
|
||||
self._warn_element_once(
|
||||
plugin_id, "redraw_vegas_element(%r) returned %s, expected an "
|
||||
"image of %dx%d; ignoring it", key,
|
||||
f"{image.width}x{image.height}" if isinstance(image, Image.Image)
|
||||
else type(image).__name__, inner, height)
|
||||
return None
|
||||
_pinned, pixels = pin_element(image, padding)
|
||||
return RenderedElement(key=key, epoch=epoch, version=None, pixels=pixels,
|
||||
digest=pixel_digest(pixels), width=pixels.shape[1])
|
||||
|
||||
def _valid_elements(self, result: Any, plugin_id: str) -> Optional[List[VegasElement]]:
|
||||
"""The usable elements in a get_vegas_elements() answer, in order.
|
||||
|
||||
None for no answer (the plugin wants its ordinary content). Anything
|
||||
that is not a VegasElement with a key and a non-empty image is
|
||||
dropped, and a duplicate key keeps its first element, each reported
|
||||
once.
|
||||
"""
|
||||
if result is None:
|
||||
return None
|
||||
@@ -903,11 +1000,8 @@ class PluginAdapter:
|
||||
plugin_id, "get_vegas_elements() returned %s, expected a list "
|
||||
"of VegasElement", type(result).__name__)
|
||||
return None
|
||||
|
||||
padding = self.config.content_padding if self.config.auto_trim else 0
|
||||
now = time.monotonic()
|
||||
seen = set()
|
||||
images: List[Image.Image] = []
|
||||
valid: List[VegasElement] = []
|
||||
for element in result:
|
||||
if not (isinstance(element, VegasElement)
|
||||
and isinstance(element.key, str) and element.key
|
||||
@@ -917,37 +1011,124 @@ class PluginAdapter:
|
||||
"not a VegasElement with a key and an image (%s); skipping it",
|
||||
type(element).__name__)
|
||||
continue
|
||||
if element.image.width <= 0 or element.image.height <= 0:
|
||||
self._warn_element_once(
|
||||
plugin_id, "get_vegas_elements() returned an empty image for "
|
||||
"%r; skipping it", element.key)
|
||||
continue
|
||||
if element.key in seen:
|
||||
self._warn_element_once(
|
||||
plugin_id, "get_vegas_elements() returned key %r twice; "
|
||||
"keeping the first", element.key)
|
||||
continue
|
||||
seen.add(element.key)
|
||||
valid.append(element)
|
||||
return valid
|
||||
|
||||
image = element.image
|
||||
if image.height != self.display_height:
|
||||
image = image.resize((image.width, self.display_height),
|
||||
Image.Resampling.LANCZOS)
|
||||
if image.mode != 'RGB':
|
||||
image = image.convert('RGB')
|
||||
def _element_image(self, element: VegasElement) -> Image.Image:
|
||||
"""An element's image at the display's height, in RGB."""
|
||||
image = element.image
|
||||
if image.height != self.display_height:
|
||||
image = image.resize((image.width, self.display_height),
|
||||
Image.Resampling.LANCZOS)
|
||||
if image.mode != 'RGB':
|
||||
image = image.convert('RGB')
|
||||
return image
|
||||
|
||||
if not element.live:
|
||||
# Plain content; a tag copied from a reused image must not
|
||||
# make it live by accident.
|
||||
images.append(untag(image.copy()) if meta_of(image) else image)
|
||||
continue
|
||||
def _padding(self) -> int:
|
||||
"""Black columns a live element carries each side: what trimming would leave."""
|
||||
return self.config.content_padding if self.config.auto_trim else 0
|
||||
|
||||
pinned, pixels = pin_element(image, padding)
|
||||
@staticmethod
|
||||
def _refresh_hz(element: VegasElement) -> float:
|
||||
try:
|
||||
return max(0.0, float(element.refresh_hz or 0.0))
|
||||
except (TypeError, ValueError):
|
||||
return 0.0
|
||||
|
||||
def _images_from_elements(
|
||||
self, result: Any, plugin_id: str, epoch: int
|
||||
) -> Optional[List[Image.Image]]:
|
||||
"""Turn get_vegas_elements()'s answer into images for the pipeline.
|
||||
|
||||
Live elements come out pinned (RGB, display height, content_padding
|
||||
black each side, never trimmed afterwards) and tagged with their
|
||||
ElementMeta; plain ones (``live=False``) come out as ordinary content.
|
||||
"""
|
||||
elements = self._valid_elements(result, plugin_id)
|
||||
if elements is None:
|
||||
return None
|
||||
padding = self._padding()
|
||||
now = time.monotonic()
|
||||
images: List[Image.Image] = []
|
||||
for element in elements:
|
||||
try:
|
||||
refresh_hz = max(0.0, float(element.refresh_hz or 0.0))
|
||||
except (TypeError, ValueError):
|
||||
refresh_hz = 0.0
|
||||
image = self._element_image(element)
|
||||
if not element.live:
|
||||
# Plain content; a tag copied from a reused image must not
|
||||
# make it live by accident.
|
||||
images.append(untag(image.copy()) if meta_of(image) else image)
|
||||
continue
|
||||
pinned, pixels = pin_element(image, padding)
|
||||
except Exception as exc: # pylint: disable=broad-except
|
||||
# An image Pillow cannot resize or convert (an odd mode, a
|
||||
# closed file) costs that element, not its neighbours.
|
||||
self._warn_element_once(
|
||||
plugin_id, "element %r could not be converted (%r); skipping it",
|
||||
element.key, exc)
|
||||
continue
|
||||
images.append(tag(pinned, ElementMeta(
|
||||
plugin_id=plugin_id, key=element.key, epoch=epoch,
|
||||
digest=pixel_digest(pixels), rendered_at=now,
|
||||
refresh_hz=refresh_hz, version=element.version)))
|
||||
refresh_hz=self._refresh_hz(element), version=element.version)))
|
||||
return images or None
|
||||
|
||||
def _rendered_from_elements(
|
||||
self, result: Any, plugin_id: str, epoch: int
|
||||
) -> Dict[str, RenderedElement]:
|
||||
"""RenderedElements for the live elements in a get_vegas_elements() answer.
|
||||
|
||||
An element handed back as the very image last converted for its key,
|
||||
with the same ``version``, reuses that conversion's pinned pixels and
|
||||
digest, so nothing is converted or checksummed. The image must be the
|
||||
same object: a plugin redrawn for a new config (new colours, a
|
||||
different font) can keep its data version, and must not keep its old
|
||||
pixels with it.
|
||||
"""
|
||||
elements = self._valid_elements(result, plugin_id) or []
|
||||
padding = self._padding()
|
||||
memo = self._element_memo
|
||||
rendered: Dict[str, RenderedElement] = {}
|
||||
for element in elements:
|
||||
if not element.live:
|
||||
continue
|
||||
memo_key = (plugin_id, element.key)
|
||||
cached = memo.get(memo_key)
|
||||
if cached is not None and element.version is not None \
|
||||
and cached[0] == element.version and cached[1] is element.image \
|
||||
and cached[2] == (padding, self.display_height):
|
||||
pixels, digest = cached[3], cached[4]
|
||||
else:
|
||||
try:
|
||||
_pinned, pixels = pin_element(self._element_image(element), padding)
|
||||
except Exception as exc: # pylint: disable=broad-except
|
||||
self._warn_element_once(
|
||||
plugin_id, "element %r could not be converted (%r); it "
|
||||
"keeps what it shows", element.key, exc)
|
||||
continue
|
||||
digest = pixel_digest(pixels)
|
||||
memo[memo_key] = (element.version, element.image,
|
||||
(padding, self.display_height), pixels, digest)
|
||||
rendered[element.key] = RenderedElement(
|
||||
key=element.key, epoch=epoch, version=element.version,
|
||||
pixels=pixels, digest=digest, width=pixels.shape[1])
|
||||
# Forget keys the plugin no longer has, so the memo stays its size.
|
||||
stale = [k for k in list(memo)
|
||||
if k[0] == plugin_id and k[1] not in rendered]
|
||||
for memo_key in stale:
|
||||
memo.pop(memo_key, None)
|
||||
return rendered
|
||||
|
||||
def _get_native_content(
|
||||
self, plugin: 'BasePlugin', plugin_id: str, restricted: bool = False
|
||||
) -> Optional[List[Image.Image]]:
|
||||
@@ -1510,6 +1691,9 @@ class PluginAdapter:
|
||||
self._content_cache.pop(plugin_id, None)
|
||||
else:
|
||||
self._content_cache.clear()
|
||||
# A config change, most often. Swapped rather than cleared:
|
||||
# the live-element worker may be iterating the old one.
|
||||
self._element_memo = {}
|
||||
|
||||
def invalidate_plugin_scroll_cache(
|
||||
self, plugin: 'BasePlugin', plugin_id: str
|
||||
|
||||
@@ -19,7 +19,8 @@ from src.common.scroll_config import solve_crisp
|
||||
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.elements import ElementMeta, ElementRecord, meta_of
|
||||
from src.vegas_mode.elements import ElementMeta, ElementRecord, LivePatch, LiveView, meta_of
|
||||
from src.vegas_mode.live_worker import LEGACY_PREFETCH_JOIN_S, VegasWorker
|
||||
from src.vegas_mode.geometry import separation_gap
|
||||
from src.vegas_mode.stream_manager import StreamManager
|
||||
|
||||
@@ -78,6 +79,26 @@ class RenderPipeline:
|
||||
# anything computed against the old strip can tell.
|
||||
_strip_gen: int = 0
|
||||
|
||||
# Live updates (see apply_live_patches and src/vegas_mode/live_worker.py).
|
||||
# Where the viewport is, published every frame for the worker.
|
||||
_view: Optional[LiveView] = None
|
||||
# Set by the coordinator for a run in which live elements are on.
|
||||
_live_enabled: bool = False
|
||||
_live_worker: Optional[VegasWorker] = None
|
||||
# The last worker stopped: it finishes its current job, and hands over
|
||||
# what it has of a group, before the one-shot prefetch fetches another.
|
||||
_retired_worker: Optional[VegasWorker] = None
|
||||
#: Patches applied between two frames at most, and the bytes they may
|
||||
#: copy, in screens of pixels (at least one patch is always applied).
|
||||
LIVE_PATCHES_PER_FRAME = 4
|
||||
LIVE_PATCH_BUDGET_SCREENS = 2
|
||||
#: A worker that dies this often in this many seconds is not restarted
|
||||
#: again this run; live updates stop and the one-shot prefetch returns.
|
||||
LIVE_WORKER_MAX_DEATHS = 3
|
||||
LIVE_WORKER_DEATH_WINDOW = 600.0
|
||||
#: Frames between checks that the worker is still alive.
|
||||
LIVE_SUPERVISE_FRAMES = 256
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
config: VegasModeConfig,
|
||||
@@ -129,6 +150,16 @@ class RenderPipeline:
|
||||
self._cycle_complete = False
|
||||
self._segments_in_scroll: List[str] = [] # Plugin IDs in current scroll
|
||||
self._record_by_seq: Dict[int, ElementRecord] = {}
|
||||
# Live updates. _applied: per record, the (epoch, digest) of the
|
||||
# pixels the strip holds. _live_slots / _live_ready: the worker's
|
||||
# hand-over, one slot per record (latest wins) and the order they
|
||||
# arrived in. Written by the worker, consumed by the render thread;
|
||||
# single-key dict operations and deque append/popleft only.
|
||||
self._applied: Dict[int, Tuple[int, Any]] = {}
|
||||
self._live_slots: Dict[int, LivePatch] = {}
|
||||
self._live_ready: Deque[int] = deque()
|
||||
self._worker_deaths: Deque[float] = deque()
|
||||
self._live_frames = 0
|
||||
|
||||
# The sub-pixel path's pacing; the crisp path solves its own (frame_interval).
|
||||
self._frame_interval = config.get_frame_interval()
|
||||
@@ -474,23 +505,43 @@ class RenderPipeline:
|
||||
if not self.config.continuous_scroll:
|
||||
return
|
||||
|
||||
# With live elements in the strip, the live-element worker fetches
|
||||
# groups too, one plugin at a time between its redraws, so that only
|
||||
# one thread ever draws for the strip.
|
||||
worker = self._live_worker
|
||||
if worker is not None and worker.is_alive():
|
||||
with self._prefetch_lock:
|
||||
if self._prepared_group is not None:
|
||||
return
|
||||
worker.request_group()
|
||||
return
|
||||
|
||||
with self._prefetch_lock:
|
||||
if self._prefetch_thread is not None and self._prefetch_thread.is_alive():
|
||||
return
|
||||
if self._prepared_group is not None:
|
||||
return # already have one waiting
|
||||
generation = self._prefetch_generation
|
||||
retired = self._retired_worker
|
||||
|
||||
def _work():
|
||||
# Deprioritise against the render loop. Linux applies nice
|
||||
# per-thread, and the heavy lifting here is PIL and numpy work
|
||||
# that releases the GIL, so the scheduler can actually act on
|
||||
# it — without this the prefetch competes for the same cores and
|
||||
# costs frames.
|
||||
# Deprioritise against the render loop for CPU time (Linux
|
||||
# applies nice per thread). Nice does nothing about the GIL,
|
||||
# which Pillow's drawing holds (docs/OFFSCREEN_RENDERING.md,
|
||||
# risk 5); the render gate below is what keeps this thread off
|
||||
# it while the render thread needs it.
|
||||
try:
|
||||
os.nice(10)
|
||||
except (OSError, AttributeError):
|
||||
pass
|
||||
# A live-element worker just stopped finishes its current job
|
||||
# and hands over what it has of a group. Wait for it, so only
|
||||
# one thread draws and a group it handed over is not replaced.
|
||||
if retired is not None and retired is not threading.current_thread() and retired.is_alive():
|
||||
retired.join(LEGACY_PREFETCH_JOIN_S)
|
||||
with self._prefetch_lock:
|
||||
if generation != self._prefetch_generation or self._prepared_group is not None:
|
||||
return
|
||||
# With vegas_scroll.prefetch_gate on, run only while the render
|
||||
# thread waits on vsync; see src/common/render_gate.py.
|
||||
gate = getattr(self.display_manager, 'render_gate', None)
|
||||
@@ -503,7 +554,8 @@ class RenderPipeline:
|
||||
with self._prefetch_lock:
|
||||
if generation != self._prefetch_generation:
|
||||
return # Vegas was reset while this was fetching
|
||||
self._prepared_group = group
|
||||
if self._prepared_group is None:
|
||||
self._prepared_group = group
|
||||
|
||||
self._prefetch_thread = threading.Thread(
|
||||
target=_work, daemon=True, name="vegas-strip-prefetch")
|
||||
@@ -825,8 +877,12 @@ class RenderPipeline:
|
||||
if new:
|
||||
self._elements = self._elements + tuple(new)
|
||||
by_seq = self.__dict__.setdefault('_record_by_seq', {})
|
||||
applied = self.__dict__.setdefault('_applied', {})
|
||||
for record in new:
|
||||
by_seq[record.seq] = record
|
||||
applied[record.seq] = (record.epoch, record.digest)
|
||||
if self._live_enabled:
|
||||
self._ensure_live_worker()
|
||||
return len(new)
|
||||
|
||||
def _forget_trimmed_records(self, cut: int) -> None:
|
||||
@@ -841,9 +897,13 @@ class RenderPipeline:
|
||||
kept = tuple(r for r in records if r.abs_x + r.width > origin)
|
||||
if len(kept) != len(records):
|
||||
by_seq = self.__dict__.setdefault('_record_by_seq', {})
|
||||
applied = self.__dict__.setdefault('_applied', {})
|
||||
slots = self.__dict__.setdefault('_live_slots', {})
|
||||
for record in records:
|
||||
if record.abs_x + record.width <= origin:
|
||||
by_seq.pop(record.seq, None)
|
||||
applied.pop(record.seq, None)
|
||||
slots.pop(record.seq, None)
|
||||
self._elements = kept
|
||||
|
||||
def _reset_records(self) -> None:
|
||||
@@ -852,6 +912,145 @@ class RenderPipeline:
|
||||
self._strip_origin = 0
|
||||
self._elements = ()
|
||||
self._record_by_seq = {}
|
||||
self._applied = {}
|
||||
# Patches still queued belong to the old strip; apply would drop them
|
||||
# on their generation anyway, but there is no reason to keep them.
|
||||
self._live_slots = {}
|
||||
self._live_ready = deque()
|
||||
self._view = None
|
||||
|
||||
def has_live_records(self) -> bool:
|
||||
"""Whether the strip holds any live element."""
|
||||
return bool(self._elements)
|
||||
|
||||
# -- live updates ---------------------------------------------------------
|
||||
|
||||
def set_live(self, enabled: bool) -> None:
|
||||
"""Switch live updates on or off for this run (the coordinator decides)."""
|
||||
self._live_enabled = enabled
|
||||
if enabled:
|
||||
if self._elements:
|
||||
self._ensure_live_worker()
|
||||
elif self._stop_live_worker():
|
||||
# The worker was fetching the strip's groups as well. Hand that
|
||||
# back to the one-shot prefetch now: nothing else asks for a group
|
||||
# until the next extension, which would find none prepared and
|
||||
# fetch inline, stalling the scroll.
|
||||
self.start_prefetch()
|
||||
|
||||
def notify_live_data(self, plugin_id: str) -> None:
|
||||
"""A plugin's data may have changed: wake the worker, if one runs."""
|
||||
worker = self._live_worker
|
||||
if worker is not None:
|
||||
worker.notify_data(plugin_id)
|
||||
|
||||
def _ensure_live_worker(self) -> None:
|
||||
"""Start the live-element worker, or restart one that died.
|
||||
|
||||
Started lazily, by the first live element placed: an install with no
|
||||
plugin that has live elements keeps the one-shot prefetch thread and
|
||||
never runs this worker at all. A worker that keeps dying is given up
|
||||
on for the run; live updates stop and the one-shot prefetch returns.
|
||||
"""
|
||||
if not self._live_enabled:
|
||||
return
|
||||
worker = self._live_worker
|
||||
if worker is not None and worker.is_alive():
|
||||
return
|
||||
if worker is None and not self._elements:
|
||||
# Nothing live in the strip: the one-shot prefetch does the work
|
||||
# until a live element is placed (_register_elements).
|
||||
return
|
||||
deaths = self.__dict__.setdefault('_worker_deaths', deque())
|
||||
now = time.monotonic()
|
||||
if worker is not None:
|
||||
deaths.append(now)
|
||||
while deaths and now - deaths[0] > self.LIVE_WORKER_DEATH_WINDOW:
|
||||
deaths.popleft()
|
||||
if len(deaths) >= self.LIVE_WORKER_MAX_DEATHS:
|
||||
logger.error(
|
||||
"Vegas live worker stopped %d times in %.0fs; live updates "
|
||||
"are off until Vegas restarts", len(deaths),
|
||||
self.LIVE_WORKER_DEATH_WINDOW)
|
||||
self._live_enabled = False
|
||||
self._live_worker = None
|
||||
# Whatever group the dead worker was fetching is lost.
|
||||
self.start_prefetch()
|
||||
return
|
||||
logger.warning("Vegas live worker was not running; restarting it")
|
||||
worker = VegasWorker(self)
|
||||
self._live_worker = worker
|
||||
worker.start()
|
||||
# Whatever the one-shot prefetch was asked for, the worker now does.
|
||||
with self._prefetch_lock:
|
||||
wanted = self._prepared_group is None
|
||||
if wanted and self.config.continuous_scroll:
|
||||
worker.request_group()
|
||||
|
||||
def _stop_live_worker(self) -> bool:
|
||||
"""Ask the worker to stop after its current job. Whether one was running."""
|
||||
worker, self._live_worker = self._live_worker, None
|
||||
if worker is None:
|
||||
return False
|
||||
self._retired_worker = worker
|
||||
worker.stop()
|
||||
return True
|
||||
|
||||
def apply_live_patches(self) -> int:
|
||||
"""Copy the worker's finished redraws into the strip. Render thread only.
|
||||
|
||||
Called between two frames (coordinator.run_frame). The only work here
|
||||
is popping prepared patches and a numpy slice copy per patch -- no
|
||||
drawing, no locks, no allocation -- bounded to LIVE_PATCHES_PER_FRAME
|
||||
patches or LIVE_PATCH_BUDGET_SCREENS screens of bytes, whichever comes
|
||||
first (always at least one). A patch is dropped when it no longer
|
||||
fits: made for an older strip, for an element trimmed away or already
|
||||
behind the screen, or older than what the strip already shows.
|
||||
|
||||
Returns:
|
||||
Patches applied.
|
||||
"""
|
||||
if self._live_enabled:
|
||||
self._live_frames = self.__dict__.get('_live_frames', 0) + 1
|
||||
if self._live_frames % self.LIVE_SUPERVISE_FRAMES == 0:
|
||||
self._ensure_live_worker()
|
||||
ready = self.__dict__.get('_live_ready')
|
||||
if not ready:
|
||||
return 0
|
||||
slots = self._live_slots
|
||||
if getattr(self, 'sync_manager', None) is not None:
|
||||
# Defensive: live elements are never on under sync, and the
|
||||
# follower would not see a patch.
|
||||
ready.clear()
|
||||
slots.clear()
|
||||
return 0
|
||||
budget = (self.LIVE_PATCH_BUDGET_SCREENS * self.display_width
|
||||
* self.display_height * 3)
|
||||
helper = self.scroll_helper
|
||||
left_edge = int(helper.scroll_position)
|
||||
applied = 0
|
||||
moved = 0
|
||||
while ready and applied < self.LIVE_PATCHES_PER_FRAME \
|
||||
and (applied == 0 or moved < budget):
|
||||
seq = ready.popleft()
|
||||
patch = slots.pop(seq, None)
|
||||
if patch is None:
|
||||
continue # a newer patch for this record already went
|
||||
record = self._record_by_seq.get(seq)
|
||||
if record is None or patch.strip_gen != self._strip_gen:
|
||||
continue
|
||||
previous = self._applied.get(seq)
|
||||
if previous is not None and patch.epoch < previous[0]:
|
||||
continue
|
||||
x = record.abs_x - self._strip_origin
|
||||
if x + record.width <= left_edge:
|
||||
continue # scrolled past; nobody will see it
|
||||
moved += helper.patch_columns(x, patch.pixels)
|
||||
self._applied[seq] = (patch.epoch, patch.digest)
|
||||
applied += 1
|
||||
if applied:
|
||||
self._note_op('patch', moved)
|
||||
return applied
|
||||
|
||||
def live_records(self) -> Tuple[ElementRecord, ...]:
|
||||
"""The live elements in the strip, in the order they were placed."""
|
||||
@@ -875,6 +1074,13 @@ class RenderPipeline:
|
||||
|
||||
# Update scroll position
|
||||
self.scroll_helper.update_scroll_position()
|
||||
# Where the viewport is now, for the live-element worker: one
|
||||
# tuple store, read by the worker without a lock.
|
||||
left = self._strip_origin + int(self.scroll_helper.scroll_position)
|
||||
self._view = LiveView(
|
||||
abs_left=left, abs_right=left + self.display_width,
|
||||
abs_end=self._strip_origin + self.scroll_helper.total_scroll_width,
|
||||
t_mono=time.monotonic())
|
||||
|
||||
# Determine if the cycle is done.
|
||||
#
|
||||
@@ -1187,6 +1393,7 @@ class RenderPipeline:
|
||||
self._prepared_group = None
|
||||
self._deferred_queue = []
|
||||
self._static_markers = ()
|
||||
self._stop_live_worker()
|
||||
self._reset_records()
|
||||
|
||||
self.display_manager.set_scrolling_state(False)
|
||||
|
||||
@@ -717,6 +717,20 @@ class StreamManager:
|
||||
A STATIC plugin is also returned with an empty list, unfetched: it
|
||||
pauses the scroll instead of adding to it (see is_static_plugin).
|
||||
"""
|
||||
group: List[Tuple[str, Optional[List[Image.Image]]]] = []
|
||||
for plugin_id in self.plan_next_group(count):
|
||||
member = self.fetch_group_member(plugin_id, offscreen_only=offscreen_only)
|
||||
if member is not None:
|
||||
group.append(member)
|
||||
return group
|
||||
|
||||
def plan_next_group(self, count: Optional[int] = None) -> List[str]:
|
||||
"""Which plugins the next group holds, advancing the rotation past them.
|
||||
|
||||
The first half of take_next_group(). The live-element worker fetches
|
||||
a group one plugin at a time (fetch_group_member) so it can fit more
|
||||
urgent redraws between them.
|
||||
"""
|
||||
if count is None:
|
||||
count = self.config.plugins_per_cycle
|
||||
|
||||
@@ -730,37 +744,38 @@ class StreamManager:
|
||||
for _ in range(min(max(1, count), total)):
|
||||
ids.append(self._ordered_plugins[self._prefetch_index])
|
||||
self._prefetch_index = (self._prefetch_index + 1) % total
|
||||
return ids
|
||||
|
||||
plugins = getattr(self.plugin_manager, 'plugins', {})
|
||||
group: List[Tuple[str, Optional[List[Image.Image]]]] = []
|
||||
def fetch_group_member(
|
||||
self, plugin_id: str, offscreen_only: bool = False
|
||||
) -> Optional[Tuple[str, Optional[List[Image.Image]]]]:
|
||||
"""One plugin's entry in a group, as take_next_group() describes it.
|
||||
|
||||
None when the plugin is gone or its fetch raised: it is left out of
|
||||
the group.
|
||||
"""
|
||||
plugin = getattr(self.plugin_manager, 'plugins', {}).get(plugin_id)
|
||||
if not plugin:
|
||||
return None
|
||||
if self.is_static_plugin(plugin_id):
|
||||
# A STATIC plugin pauses the scroll rather than scrolling by, so it
|
||||
# contributes no columns. It keeps its place in the group (empty)
|
||||
# so the pipeline can mark where its turn falls.
|
||||
return (plugin_id, [])
|
||||
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
|
||||
return None
|
||||
if images:
|
||||
self.stats['segments_fetched'] += 1
|
||||
return (plugin_id, images)
|
||||
# Only the old contract hands anything back to the render thread.
|
||||
defer_empty = offscreen_only and not getattr(
|
||||
self.config, 'offscreen_prefetch', True)
|
||||
|
||||
for plugin_id in ids:
|
||||
plugin = plugins.get(plugin_id)
|
||||
if not plugin:
|
||||
continue
|
||||
if self.is_static_plugin(plugin_id):
|
||||
# A STATIC plugin pauses the scroll rather than scrolling by,
|
||||
# so it contributes no columns. It keeps its place in the
|
||||
# group (empty) so the pipeline can mark where its turn falls.
|
||||
group.append((plugin_id, []))
|
||||
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))
|
||||
else:
|
||||
group.append((plugin_id, None if defer_empty else []))
|
||||
|
||||
return group
|
||||
return (plugin_id, None if defer_empty else [])
|
||||
|
||||
def advance_cycle(self) -> None:
|
||||
"""
|
||||
|
||||
Reference in New Issue
Block a user