mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-10 17:16:36 +00:00
experiment(vegas): vegas_scroll.prefetch_gate runs the prefetch only while the render thread waits on vsync
The render thread spends most of each refresh in SwapOnVSync with the GIL released, then needs it back the moment the swap returns. With plugin rendering on the prefetch thread, it often has to wait for it -- behind bytecode for up to the switch interval, behind a GIL-holding C call for as long as that takes -- and hdpi's late frames of 2-5 refreshes went up. src/common/render_gate.py opens a window around each swap, up to just before the refresh the swap will return on, and a profile hook on the prefetch thread parks it outside that window. It is never parked holding a lock the render thread also takes (the Vegas buffer, cache and state locks, logging, threading, importlib, the cache), never when no frame has been swapped for 50ms, and never for more than 50ms at a time. Off by default and ignored on a binding that keeps the GIL in SwapOnVSync, where the window would never let the prefetch run. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,222 @@
|
|||||||
|
"""Let a background thread run Python only while the render thread waits on vsync.
|
||||||
|
|
||||||
|
With plugin rendering moved to Vegas's prefetch thread (DisplayManager.offscreen,
|
||||||
|
#630) the render thread no longer stops for it, but it still shares the GIL
|
||||||
|
with it. The render thread spends most of each refresh inside SwapOnVSync,
|
||||||
|
which releases the GIL, and needs it back the moment the swap returns. If the
|
||||||
|
prefetch thread is running Python right then, the render thread waits: up to
|
||||||
|
the switch interval (5ms) behind bytecode, and for as long as a C call that
|
||||||
|
keeps the GIL takes. On hdpi that showed up as frames 2-5 refreshes late while
|
||||||
|
a group was being prepared.
|
||||||
|
|
||||||
|
The gate turns that around. The display manager opens it just before each swap,
|
||||||
|
with a deadline shortly ahead of the refresh the swap will return on, and
|
||||||
|
closes it when the swap returns. A thread inside ``gate.yielding()`` checks it on
|
||||||
|
every Python and C call through a profile hook, and once the window has closed
|
||||||
|
it parks -- blocked on a condition, GIL released -- until the next swap opens
|
||||||
|
it. The render thread then finds the GIL free when its refresh arrives, and the
|
||||||
|
background work runs in time the render thread was only spending waiting.
|
||||||
|
|
||||||
|
Parking a thread is only safe if nothing the render thread needs is stuck
|
||||||
|
behind it, so it is never parked:
|
||||||
|
|
||||||
|
* while it holds a lock registered with ``guard()`` (the Vegas buffers and
|
||||||
|
caches the render thread also takes);
|
||||||
|
* inside logging, threading, importlib or the cache, all of which take locks the
|
||||||
|
render thread can take too;
|
||||||
|
* when there is no render loop to protect -- no swap for ``STALE_SECONDS``, as
|
||||||
|
on a static screen or a stalled frame.
|
||||||
|
|
||||||
|
And a parked thread is never held more than ``MAX_WAIT_SECONDS`` at a time, so
|
||||||
|
whatever the gate gets wrong costs a frame, not a freeze.
|
||||||
|
|
||||||
|
Only worth enabling on a binding whose SwapOnVSync releases the GIL (see
|
||||||
|
scripts/build_rgbmatrix_nogil.sh). With one that keeps it, the window never
|
||||||
|
lets the background thread run and it only makes progress in the timeouts.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import math
|
||||||
|
import sys
|
||||||
|
import threading
|
||||||
|
import time
|
||||||
|
from collections import deque
|
||||||
|
from typing import Any, Callable, Deque, List, Optional
|
||||||
|
|
||||||
|
#: Park background threads this long before the refresh a swap will return on,
|
||||||
|
#: so a short C call already under way has finished by then.
|
||||||
|
MARGIN_SECONDS = 0.002
|
||||||
|
|
||||||
|
#: The longest a background thread is parked in one go.
|
||||||
|
MAX_WAIT_SECONDS = 0.05
|
||||||
|
|
||||||
|
#: No swap for this long means there is no render loop running to protect.
|
||||||
|
STALE_SECONDS = 0.05
|
||||||
|
|
||||||
|
#: Swaps needed before the refresh period is trusted enough to open a window.
|
||||||
|
MIN_SAMPLES = 8
|
||||||
|
|
||||||
|
#: Parking with any of these on the stack could hold a lock the render
|
||||||
|
#: thread takes: logging handler locks, Condition and Event internals, the
|
||||||
|
#: module import locks, and the disk and memory cache locks.
|
||||||
|
_UNSAFE_PATHS = (
|
||||||
|
"logging",
|
||||||
|
"threading.py",
|
||||||
|
"importlib",
|
||||||
|
"cache",
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _unsafe(frame: Any, base: Any) -> bool:
|
||||||
|
"""True if a frame above ``base`` comes from somewhere parking could deadlock.
|
||||||
|
|
||||||
|
``base`` is the frame that entered ``yielding()``; what lies below it (the
|
||||||
|
thread's own bootstrap in threading.py) holds nothing.
|
||||||
|
"""
|
||||||
|
while frame is not None and frame is not base:
|
||||||
|
filename = frame.f_code.co_filename
|
||||||
|
for part in _UNSAFE_PATHS:
|
||||||
|
if part in filename:
|
||||||
|
return True
|
||||||
|
frame = frame.f_back
|
||||||
|
return False
|
||||||
|
|
||||||
|
|
||||||
|
def swap_releases_gil() -> Optional[bool]:
|
||||||
|
"""Whether the loaded rgbmatrix binding releases the GIL, or None if none is loaded.
|
||||||
|
|
||||||
|
The rebuilt binding links PyEval_SaveThread and the stock one never does.
|
||||||
|
The same test as src.common.frame_timing.binding_releases_gil (#629); one
|
||||||
|
of the two goes once both have landed.
|
||||||
|
"""
|
||||||
|
module = sys.modules.get("rgbmatrix.core")
|
||||||
|
path = getattr(module, "__file__", None)
|
||||||
|
if not path:
|
||||||
|
return None
|
||||||
|
try:
|
||||||
|
with open(path, "rb") as handle:
|
||||||
|
return b"PyEval_SaveThread" in handle.read()
|
||||||
|
except OSError:
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
def _held(lock: Any) -> bool:
|
||||||
|
"""Is ``lock`` held? RLocks report this thread's ownership; plain locks, anyone's."""
|
||||||
|
is_owned = getattr(lock, "_is_owned", None)
|
||||||
|
if is_owned is not None:
|
||||||
|
return is_owned()
|
||||||
|
return lock.locked()
|
||||||
|
|
||||||
|
|
||||||
|
class RenderGate:
|
||||||
|
"""Opened by the render thread around each swap; honoured by background threads."""
|
||||||
|
|
||||||
|
def __init__(self, clock: Callable[[], float] = time.monotonic):
|
||||||
|
self.clock = clock
|
||||||
|
self._cond = threading.Condition()
|
||||||
|
self._generation = 0
|
||||||
|
self._open_until = 0.0
|
||||||
|
self._last_return: Optional[float] = None
|
||||||
|
self._periods: Deque[float] = deque(maxlen=64)
|
||||||
|
self._period: Optional[float] = None
|
||||||
|
self._guarded: List[Any] = []
|
||||||
|
self._local = threading.local()
|
||||||
|
#: How often, and for how long in all, background threads were parked.
|
||||||
|
self.parks = 0
|
||||||
|
self.parked_seconds = 0.0
|
||||||
|
|
||||||
|
def guard(self, *locks: Any) -> None:
|
||||||
|
"""Never park a thread while it holds (or, for a plain Lock, anyone holds) these."""
|
||||||
|
self._guarded.extend(lock for lock in locks if lock is not None)
|
||||||
|
|
||||||
|
# -- render thread -----------------------------------------------------
|
||||||
|
|
||||||
|
def refresh_period(self) -> Optional[float]:
|
||||||
|
"""The panel's refresh period from recent swaps, or None until known.
|
||||||
|
|
||||||
|
The 10th percentile of the gaps between swap returns, each divided by
|
||||||
|
the hold: a late frame only ever lengthens a gap, so the low end is
|
||||||
|
the panel's own period.
|
||||||
|
"""
|
||||||
|
return self._period
|
||||||
|
|
||||||
|
def before_swap(self, hold: int) -> None:
|
||||||
|
"""The render thread is about to block in SwapOnVSync: open the window."""
|
||||||
|
hold = max(1, int(hold))
|
||||||
|
period = self._period
|
||||||
|
now = self.clock()
|
||||||
|
last = self._last_return
|
||||||
|
if period and last is not None and now - last < STALE_SECONDS:
|
||||||
|
# The swap returns on the first refresh boundary after both the
|
||||||
|
# current frame's hold is up and this frame has been handed over;
|
||||||
|
# boundaries fall a whole period apart from the last return.
|
||||||
|
refreshes = max(hold, math.ceil((now - last) / period))
|
||||||
|
open_until = last + refreshes * period - MARGIN_SECONDS
|
||||||
|
else:
|
||||||
|
open_until = 0.0 # no rhythm to predict from: leave threads be
|
||||||
|
with self._cond:
|
||||||
|
self._open_until = open_until
|
||||||
|
self._generation += 1
|
||||||
|
self._cond.notify_all()
|
||||||
|
|
||||||
|
def after_swap(self, hold: int) -> None:
|
||||||
|
"""The swap returned and the render thread needs the GIL: close the window."""
|
||||||
|
now = self.clock()
|
||||||
|
self._open_until = 0.0
|
||||||
|
last = self._last_return
|
||||||
|
if last is not None and now - last < STALE_SECONDS:
|
||||||
|
self._periods.append((now - last) / max(1, int(hold)))
|
||||||
|
if len(self._periods) >= MIN_SAMPLES:
|
||||||
|
ordered = sorted(self._periods)
|
||||||
|
self._period = ordered[len(ordered) // 10]
|
||||||
|
self._last_return = now
|
||||||
|
|
||||||
|
# -- background threads ------------------------------------------------
|
||||||
|
|
||||||
|
def _should_park(self, frame: Any, now: float) -> bool:
|
||||||
|
if now < self._open_until:
|
||||||
|
return False # inside the window
|
||||||
|
last = self._last_return
|
||||||
|
if last is None or now - last > STALE_SECONDS or self._period is None:
|
||||||
|
return False # no render loop to protect
|
||||||
|
for lock in self._guarded:
|
||||||
|
if _held(lock):
|
||||||
|
return False
|
||||||
|
return not _unsafe(frame, getattr(self._local, "base", None))
|
||||||
|
|
||||||
|
def _hook(self, frame: Any, _event: str, _arg: Any) -> None:
|
||||||
|
now = self.clock()
|
||||||
|
if not self._should_park(frame, now):
|
||||||
|
return
|
||||||
|
generation = self._generation
|
||||||
|
with self._cond:
|
||||||
|
self._cond.wait_for(lambda: self._generation != generation,
|
||||||
|
timeout=MAX_WAIT_SECONDS)
|
||||||
|
self.parks += 1
|
||||||
|
self.parked_seconds += self.clock() - now
|
||||||
|
|
||||||
|
def yielding(self) -> "_Yielding":
|
||||||
|
"""``with gate.yielding():`` runs the block giving way to the render thread."""
|
||||||
|
return _Yielding(self)
|
||||||
|
|
||||||
|
|
||||||
|
class _Yielding:
|
||||||
|
"""Installs a gate's profile hook on the thread for the length of a block."""
|
||||||
|
|
||||||
|
def __init__(self, gate: RenderGate):
|
||||||
|
self.gate = gate
|
||||||
|
self._previous: Any = None
|
||||||
|
self._previous_base: Any = None
|
||||||
|
|
||||||
|
def __enter__(self) -> RenderGate:
|
||||||
|
local = self.gate._local # pylint: disable=protected-access
|
||||||
|
self._previous_base = getattr(local, "base", None)
|
||||||
|
local.base = sys._getframe(1) # pylint: disable=protected-access
|
||||||
|
self._previous = sys.getprofile()
|
||||||
|
sys.setprofile(self.gate._hook) # pylint: disable=protected-access
|
||||||
|
return self.gate
|
||||||
|
|
||||||
|
def __exit__(self, *_exc: Any) -> None:
|
||||||
|
sys.setprofile(self._previous)
|
||||||
|
self.gate._local.base = self._previous_base # pylint: disable=protected-access
|
||||||
@@ -322,6 +322,11 @@ class DisplayManager:
|
|||||||
# See src/common/scroll_config.py and scripts/scroll_speeds.py.
|
# See src/common/scroll_config.py and scripts/scroll_speeds.py.
|
||||||
self._frame_hold = 1
|
self._frame_hold = 1
|
||||||
|
|
||||||
|
# A src.common.render_gate.RenderGate while Vegas runs with
|
||||||
|
# vegas_scroll.prefetch_gate on: opened around each swap so the
|
||||||
|
# prefetch thread only runs Python while this thread waits on vsync.
|
||||||
|
self.render_gate = None
|
||||||
|
|
||||||
self._scrolling_state = {
|
self._scrolling_state = {
|
||||||
'is_scrolling': False,
|
'is_scrolling': False,
|
||||||
'last_scroll_activity': 0,
|
'last_scroll_activity': 0,
|
||||||
@@ -948,7 +953,12 @@ class DisplayManager:
|
|||||||
# Swap buffers immediately. framerate_fraction holds the frame
|
# Swap buffers immediately. framerate_fraction holds the frame
|
||||||
# for N refreshes; SwapOnVSync blocks for all of them, which is
|
# for N refreshes; SwapOnVSync blocks for all of them, which is
|
||||||
# what paces the render loop to the chosen frame rate.
|
# what paces the render loop to the chosen frame rate.
|
||||||
|
gate = self.render_gate
|
||||||
|
if gate is not None:
|
||||||
|
gate.before_swap(self._frame_hold)
|
||||||
self.matrix.SwapOnVSync(self.offscreen_canvas, self._frame_hold)
|
self.matrix.SwapOnVSync(self.offscreen_canvas, self._frame_hold)
|
||||||
|
if gate is not None:
|
||||||
|
gate.after_swap(self._frame_hold)
|
||||||
|
|
||||||
# Swap our canvas references
|
# Swap our canvas references
|
||||||
self.offscreen_canvas, self.current_canvas = self.current_canvas, self.offscreen_canvas
|
self.offscreen_canvas, self.current_canvas = self.current_canvas, self.offscreen_canvas
|
||||||
|
|||||||
@@ -90,6 +90,14 @@ class VegasModeConfig:
|
|||||||
# Experimental: measured with scripts/frame_soak.py before it gets a default.
|
# Experimental: measured with scripts/frame_soak.py before it gets a default.
|
||||||
switch_interval_ms: float = 0.0
|
switch_interval_ms: float = 0.0
|
||||||
|
|
||||||
|
# Let the prefetch thread run Python only while the render thread is
|
||||||
|
# blocked waiting for vsync, and park it the rest of the time, so the
|
||||||
|
# render thread never waits for the GIL when its refresh comes round. Needs
|
||||||
|
# a binding that releases the GIL in SwapOnVSync; ignored otherwise.
|
||||||
|
# Experimental: see src/common/render_gate.py and measure with
|
||||||
|
# scripts/frame_soak.py before it gets a default.
|
||||||
|
prefetch_gate: bool = False
|
||||||
|
|
||||||
# Keep one continuous strip, extending it with the next group of plugins as
|
# Keep one continuous strip, extending it with the next group of plugins as
|
||||||
# the scroll approaches the end, instead of composing a fresh strip and
|
# the scroll approaches the end, instead of composing a fresh strip and
|
||||||
# swapping it in. A swap stops the motion, substitutes every pixel at once
|
# swapping it in. A swap stops the motion, substitutes every pixel at once
|
||||||
@@ -221,6 +229,7 @@ class VegasModeConfig:
|
|||||||
continuous_scroll=vegas_config.get('continuous_scroll', True),
|
continuous_scroll=vegas_config.get('continuous_scroll', True),
|
||||||
offscreen_prefetch=bool(vegas_config.get('offscreen_prefetch', True)),
|
offscreen_prefetch=bool(vegas_config.get('offscreen_prefetch', True)),
|
||||||
switch_interval_ms=float(vegas_config.get('switch_interval_ms', 0.0) or 0.0),
|
switch_interval_ms=float(vegas_config.get('switch_interval_ms', 0.0) or 0.0),
|
||||||
|
prefetch_gate=bool(vegas_config.get('prefetch_gate', False)),
|
||||||
extend_threshold_screens=float(
|
extend_threshold_screens=float(
|
||||||
vegas_config.get('extend_threshold_screens', 2.0)),
|
vegas_config.get('extend_threshold_screens', 2.0)),
|
||||||
auto_trim=vegas_config.get('auto_trim', True),
|
auto_trim=vegas_config.get('auto_trim', True),
|
||||||
@@ -264,6 +273,7 @@ class VegasModeConfig:
|
|||||||
'continuous_scroll': self.continuous_scroll,
|
'continuous_scroll': self.continuous_scroll,
|
||||||
'offscreen_prefetch': self.offscreen_prefetch,
|
'offscreen_prefetch': self.offscreen_prefetch,
|
||||||
'switch_interval_ms': self.switch_interval_ms,
|
'switch_interval_ms': self.switch_interval_ms,
|
||||||
|
'prefetch_gate': self.prefetch_gate,
|
||||||
'extend_threshold_screens': self.extend_threshold_screens,
|
'extend_threshold_screens': self.extend_threshold_screens,
|
||||||
'auto_trim': self.auto_trim,
|
'auto_trim': self.auto_trim,
|
||||||
'trim_threshold': self.trim_threshold,
|
'trim_threshold': self.trim_threshold,
|
||||||
|
|||||||
@@ -18,6 +18,7 @@ import time
|
|||||||
import threading
|
import threading
|
||||||
from typing import Optional, Dict, Any, List, Callable, TYPE_CHECKING
|
from typing import Optional, Dict, Any, List, Callable, TYPE_CHECKING
|
||||||
|
|
||||||
|
from src.common import render_gate
|
||||||
from src.vegas_mode.config import VegasModeConfig
|
from src.vegas_mode.config import VegasModeConfig
|
||||||
from src.vegas_mode.plugin_adapter import PluginAdapter
|
from src.vegas_mode.plugin_adapter import PluginAdapter
|
||||||
from src.vegas_mode.stream_manager import StreamManager
|
from src.vegas_mode.stream_manager import StreamManager
|
||||||
@@ -282,6 +283,7 @@ class VegasModeCoordinator:
|
|||||||
self._fps_last_health_log = 0.0
|
self._fps_last_health_log = 0.0
|
||||||
self._fps_was_degraded = False
|
self._fps_was_degraded = False
|
||||||
self._apply_switch_interval()
|
self._apply_switch_interval()
|
||||||
|
self._install_render_gate()
|
||||||
|
|
||||||
# Line up the next group immediately, so the first extension is already
|
# Line up the next group immediately, so the first extension is already
|
||||||
# warm rather than stalling the scroll to fetch it.
|
# warm rather than stalling the scroll to fetch it.
|
||||||
@@ -305,6 +307,7 @@ class VegasModeCoordinator:
|
|||||||
self._start_time = None
|
self._start_time = None
|
||||||
|
|
||||||
self._restore_switch_interval()
|
self._restore_switch_interval()
|
||||||
|
self._remove_render_gate()
|
||||||
|
|
||||||
# Cleanup components
|
# Cleanup components
|
||||||
self.render_pipeline.reset()
|
self.render_pipeline.reset()
|
||||||
@@ -330,6 +333,38 @@ class VegasModeCoordinator:
|
|||||||
sys.setswitchinterval(saved)
|
sys.setswitchinterval(saved)
|
||||||
self._saved_switch_interval = None
|
self._saved_switch_interval = None
|
||||||
|
|
||||||
|
def _install_render_gate(self) -> None:
|
||||||
|
"""Gate the prefetch thread on the render thread's swaps; see VegasModeConfig."""
|
||||||
|
if not self.vegas_config.prefetch_gate:
|
||||||
|
return
|
||||||
|
if getattr(self.display_manager, 'render_gate', None) is not None:
|
||||||
|
return
|
||||||
|
releases = render_gate.swap_releases_gil()
|
||||||
|
if not releases:
|
||||||
|
logger.warning(
|
||||||
|
"Vegas: prefetch_gate ignored -- %s",
|
||||||
|
"this rgbmatrix binding keeps the GIL in SwapOnVSync "
|
||||||
|
"(scripts/build_rgbmatrix_nogil.sh)" if releases is False
|
||||||
|
else "no hardware binding loaded")
|
||||||
|
return
|
||||||
|
gate = render_gate.RenderGate()
|
||||||
|
# Locks the render thread takes too: never park the prefetch holding one.
|
||||||
|
gate.guard(self._state_lock,
|
||||||
|
getattr(self.stream_manager, '_buffer_lock', None),
|
||||||
|
getattr(self.render_pipeline, '_buffer_lock', None),
|
||||||
|
getattr(self.render_pipeline, '_prefetch_lock', None),
|
||||||
|
getattr(self.plugin_adapter, '_cache_lock', None))
|
||||||
|
self.display_manager.render_gate = gate
|
||||||
|
logger.info("Vegas: prefetch gated on vsync")
|
||||||
|
|
||||||
|
def _remove_render_gate(self) -> None:
|
||||||
|
gate = getattr(self.display_manager, 'render_gate', None)
|
||||||
|
if gate is None:
|
||||||
|
return
|
||||||
|
self.display_manager.render_gate = None
|
||||||
|
logger.info("Vegas: prefetch gate parked the prefetch %d times, %.1fs in all",
|
||||||
|
gate.parks, gate.parked_seconds)
|
||||||
|
|
||||||
def pause(self) -> None:
|
def pause(self) -> None:
|
||||||
"""Pause Vegas mode (for live priority interruption)."""
|
"""Pause Vegas mode (for live priority interruption)."""
|
||||||
with self._state_lock:
|
with self._state_lock:
|
||||||
@@ -552,6 +587,10 @@ class VegasModeCoordinator:
|
|||||||
fps, target, fps_frame_count,
|
fps, target, fps_frame_count,
|
||||||
p99 * 1000.0, frame_worst * 1000.0
|
p99 * 1000.0, frame_worst * 1000.0
|
||||||
)
|
)
|
||||||
|
gate = getattr(self.display_manager, 'render_gate', None)
|
||||||
|
if gate is not None:
|
||||||
|
logger.info("Vegas: prefetch parked %d times, %.1fs in all",
|
||||||
|
gate.parks, gate.parked_seconds)
|
||||||
self._fps_last_health_log = current_time
|
self._fps_last_health_log = current_time
|
||||||
else:
|
else:
|
||||||
logger.debug(
|
logger.debug(
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ import os
|
|||||||
import time
|
import time
|
||||||
import threading
|
import threading
|
||||||
from collections import deque
|
from collections import deque
|
||||||
|
from contextlib import nullcontext
|
||||||
from typing import Optional, List, Any, Dict, Deque, TYPE_CHECKING
|
from typing import Optional, List, Any, Dict, Deque, TYPE_CHECKING
|
||||||
from PIL import Image
|
from PIL import Image
|
||||||
|
|
||||||
@@ -392,7 +393,11 @@ class RenderPipeline:
|
|||||||
os.nice(10)
|
os.nice(10)
|
||||||
except (OSError, AttributeError):
|
except (OSError, AttributeError):
|
||||||
pass
|
pass
|
||||||
|
# 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)
|
||||||
try:
|
try:
|
||||||
|
with gate.yielding() if gate is not None else nullcontext():
|
||||||
group = self.stream_manager.take_next_group(offscreen_only=True)
|
group = self.stream_manager.take_next_group(offscreen_only=True)
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.exception("Background prefetch failed")
|
logger.exception("Background prefetch failed")
|
||||||
|
|||||||
@@ -0,0 +1,369 @@
|
|||||||
|
"""The render gate (src/common/render_gate.py) and its wiring into Vegas."""
|
||||||
|
|
||||||
|
import os
|
||||||
|
import sys
|
||||||
|
import threading
|
||||||
|
import time
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
os.environ.setdefault("EMULATOR", "true")
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
|
||||||
|
|
||||||
|
from src.common import render_gate # noqa: E402
|
||||||
|
from src.common.render_gate import RenderGate # noqa: E402
|
||||||
|
|
||||||
|
PERIOD = 0.010
|
||||||
|
|
||||||
|
|
||||||
|
class Clock:
|
||||||
|
def __init__(self):
|
||||||
|
self.now = 100.0
|
||||||
|
|
||||||
|
def __call__(self):
|
||||||
|
return self.now
|
||||||
|
|
||||||
|
|
||||||
|
def _running(gate, clock, swaps=render_gate.MIN_SAMPLES + 1, hold=1):
|
||||||
|
"""Drive ``swaps`` on-time swaps through the gate, a refresh apart."""
|
||||||
|
for _ in range(swaps):
|
||||||
|
gate.before_swap(hold)
|
||||||
|
clock.now += hold * PERIOD
|
||||||
|
gate.after_swap(hold)
|
||||||
|
|
||||||
|
|
||||||
|
class TestWindow:
|
||||||
|
def test_no_window_until_the_period_is_known(self):
|
||||||
|
clock = Clock()
|
||||||
|
gate = RenderGate(clock)
|
||||||
|
_running(gate, clock, swaps=3)
|
||||||
|
assert gate.refresh_period() is None
|
||||||
|
gate.before_swap(1)
|
||||||
|
assert gate._open_until == 0.0
|
||||||
|
# ...and nothing parks meanwhile.
|
||||||
|
assert not gate._should_park(sys._getframe(), clock.now + 0.009)
|
||||||
|
|
||||||
|
def test_the_period_is_the_low_end_of_the_swap_gaps(self):
|
||||||
|
clock = Clock()
|
||||||
|
gate = RenderGate(clock)
|
||||||
|
_running(gate, clock, swaps=20)
|
||||||
|
gate.before_swap(1)
|
||||||
|
clock.now += 0.030 # one late frame lengthens its gap only
|
||||||
|
gate.after_swap(1)
|
||||||
|
assert gate.refresh_period() == pytest.approx(PERIOD)
|
||||||
|
|
||||||
|
def test_opens_until_just_before_the_refresh_the_swap_returns_on(self):
|
||||||
|
clock = Clock()
|
||||||
|
gate = RenderGate(clock)
|
||||||
|
_running(gate, clock)
|
||||||
|
last = clock.now
|
||||||
|
clock.now += 0.003 # the render thread's own work
|
||||||
|
gate.before_swap(1)
|
||||||
|
assert gate._open_until == pytest.approx(
|
||||||
|
last + PERIOD - render_gate.MARGIN_SECONDS)
|
||||||
|
|
||||||
|
def test_a_held_frame_opens_for_the_whole_hold(self):
|
||||||
|
clock = Clock()
|
||||||
|
gate = RenderGate(clock)
|
||||||
|
_running(gate, clock, hold=2)
|
||||||
|
last = clock.now
|
||||||
|
clock.now += 0.003
|
||||||
|
gate.before_swap(2)
|
||||||
|
assert gate._open_until == pytest.approx(
|
||||||
|
last + 2 * PERIOD - render_gate.MARGIN_SECONDS)
|
||||||
|
|
||||||
|
def test_a_late_frame_opens_until_the_next_boundary(self):
|
||||||
|
clock = Clock()
|
||||||
|
gate = RenderGate(clock)
|
||||||
|
_running(gate, clock)
|
||||||
|
last = clock.now
|
||||||
|
clock.now += 0.013 # missed its refresh: returns on the next
|
||||||
|
gate.before_swap(1)
|
||||||
|
assert gate._open_until == pytest.approx(
|
||||||
|
last + 2 * PERIOD - render_gate.MARGIN_SECONDS)
|
||||||
|
|
||||||
|
def test_closed_the_moment_the_swap_returns(self):
|
||||||
|
clock = Clock()
|
||||||
|
gate = RenderGate(clock)
|
||||||
|
_running(gate, clock)
|
||||||
|
gate.before_swap(1)
|
||||||
|
clock.now += PERIOD
|
||||||
|
gate.after_swap(1)
|
||||||
|
assert gate._open_until == 0.0
|
||||||
|
|
||||||
|
|
||||||
|
class TestWhenToPark:
|
||||||
|
@pytest.fixture
|
||||||
|
def gate(self):
|
||||||
|
clock = Clock()
|
||||||
|
gate = RenderGate(clock)
|
||||||
|
_running(gate, clock)
|
||||||
|
gate.before_swap(1)
|
||||||
|
return gate
|
||||||
|
|
||||||
|
def test_runs_inside_the_window(self, gate):
|
||||||
|
assert not gate._should_park(sys._getframe(), gate._open_until - 0.001)
|
||||||
|
|
||||||
|
def test_parks_once_it_has_closed(self, gate):
|
||||||
|
assert gate._should_park(sys._getframe(), gate._open_until + 0.0001)
|
||||||
|
|
||||||
|
def test_never_with_no_render_loop_to_protect(self, gate):
|
||||||
|
stale = gate._last_return + render_gate.STALE_SECONDS + 0.001
|
||||||
|
assert not gate._should_park(sys._getframe(), stale)
|
||||||
|
|
||||||
|
def test_never_holding_a_guarded_lock(self, gate):
|
||||||
|
rlock, lock = threading.RLock(), threading.Lock()
|
||||||
|
gate.guard(rlock, lock, None)
|
||||||
|
closed = gate._open_until + 0.0001
|
||||||
|
with rlock:
|
||||||
|
assert not gate._should_park(sys._getframe(), closed)
|
||||||
|
with lock:
|
||||||
|
assert not gate._should_park(sys._getframe(), closed)
|
||||||
|
assert gate._should_park(sys._getframe(), closed)
|
||||||
|
|
||||||
|
def test_an_rlock_held_by_another_thread_is_not_this_ones(self, gate):
|
||||||
|
rlock = threading.RLock()
|
||||||
|
gate.guard(rlock)
|
||||||
|
taken, done = threading.Event(), threading.Event()
|
||||||
|
|
||||||
|
def hold():
|
||||||
|
with rlock:
|
||||||
|
taken.set()
|
||||||
|
done.wait(5)
|
||||||
|
other = threading.Thread(target=hold)
|
||||||
|
other.start()
|
||||||
|
try:
|
||||||
|
taken.wait(5)
|
||||||
|
assert gate._should_park(sys._getframe(), gate._open_until + 0.0001)
|
||||||
|
finally:
|
||||||
|
done.set()
|
||||||
|
other.join()
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("filename", [
|
||||||
|
"/usr/lib/python3.11/logging/__init__.py",
|
||||||
|
"/usr/lib/python3.11/threading.py",
|
||||||
|
"<frozen importlib._bootstrap>",
|
||||||
|
"/home/pi/LEDMatrix/src/cache/disk_cache.py",
|
||||||
|
"/home/pi/LEDMatrix/src/cache_manager.py",
|
||||||
|
])
|
||||||
|
def test_never_inside_code_that_takes_shared_locks(self, gate, filename):
|
||||||
|
namespace = {}
|
||||||
|
exec(compile("import sys\ndef here():\n return sys._getframe()\n",
|
||||||
|
filename, "exec"), namespace)
|
||||||
|
frame = namespace["here"]()
|
||||||
|
assert render_gate._unsafe(frame, None)
|
||||||
|
assert not gate._should_park(frame, gate._open_until + 0.0001)
|
||||||
|
|
||||||
|
def test_what_lies_below_the_yielding_block_does_not_count(self):
|
||||||
|
# A thread's stack always starts in threading.py; only frames above
|
||||||
|
# the one that entered yielding() matter.
|
||||||
|
namespace = {}
|
||||||
|
exec(compile("def bootstrap(fn):\n return fn()\n",
|
||||||
|
"/usr/lib/python3.11/threading.py", "exec"), namespace)
|
||||||
|
|
||||||
|
def entered():
|
||||||
|
base = sys._getframe()
|
||||||
|
|
||||||
|
def work():
|
||||||
|
top = sys._getframe()
|
||||||
|
return render_gate._unsafe(top, None), render_gate._unsafe(top, base)
|
||||||
|
return work()
|
||||||
|
assert namespace["bootstrap"](entered) == (True, False)
|
||||||
|
|
||||||
|
|
||||||
|
class TestYielding:
|
||||||
|
"""A real background thread, parked and released by the gate."""
|
||||||
|
|
||||||
|
def _worker(self, gate, stop):
|
||||||
|
count = [0]
|
||||||
|
|
||||||
|
def step():
|
||||||
|
count[0] += 1
|
||||||
|
|
||||||
|
def work():
|
||||||
|
with gate.yielding():
|
||||||
|
while not stop.is_set():
|
||||||
|
step()
|
||||||
|
thread = threading.Thread(target=work, daemon=True)
|
||||||
|
return thread, count
|
||||||
|
|
||||||
|
def test_parks_while_closed_and_runs_while_open(self):
|
||||||
|
clock = Clock()
|
||||||
|
gate = RenderGate(clock)
|
||||||
|
_running(gate, clock)
|
||||||
|
gate.before_swap(1)
|
||||||
|
stop = threading.Event()
|
||||||
|
thread, count = self._worker(gate, stop)
|
||||||
|
try:
|
||||||
|
clock.now = gate._open_until + 0.001 # window closed
|
||||||
|
thread.start()
|
||||||
|
time.sleep(0.2)
|
||||||
|
parked_steps = count[0]
|
||||||
|
# At most one step per MAX_WAIT timeout while parked.
|
||||||
|
assert parked_steps < 0.2 / render_gate.MAX_WAIT_SECONDS + 5
|
||||||
|
assert gate.parks >= 1
|
||||||
|
|
||||||
|
with gate._cond: # the next swap opens it
|
||||||
|
gate._open_until = clock.now + 3600.0
|
||||||
|
gate._generation += 1
|
||||||
|
gate._cond.notify_all()
|
||||||
|
time.sleep(0.1)
|
||||||
|
assert count[0] > parked_steps + 1000
|
||||||
|
finally:
|
||||||
|
stop.set()
|
||||||
|
with gate._cond:
|
||||||
|
gate._open_until = float("inf")
|
||||||
|
gate._generation += 1
|
||||||
|
gate._cond.notify_all()
|
||||||
|
thread.join(2)
|
||||||
|
assert not thread.is_alive()
|
||||||
|
|
||||||
|
def test_runs_freely_once_the_render_loop_stops(self):
|
||||||
|
clock = Clock()
|
||||||
|
gate = RenderGate(clock)
|
||||||
|
_running(gate, clock)
|
||||||
|
stop = threading.Event()
|
||||||
|
thread, count = self._worker(gate, stop)
|
||||||
|
clock.now += render_gate.STALE_SECONDS + 0.01
|
||||||
|
thread.start()
|
||||||
|
time.sleep(0.1)
|
||||||
|
stop.set()
|
||||||
|
thread.join(2)
|
||||||
|
assert count[0] > 1000
|
||||||
|
assert gate.parks == 0
|
||||||
|
|
||||||
|
def test_the_hook_comes_off_when_the_block_ends(self):
|
||||||
|
gate = RenderGate()
|
||||||
|
before = sys.getprofile()
|
||||||
|
with gate.yielding():
|
||||||
|
assert sys.getprofile() == gate._hook
|
||||||
|
assert sys.getprofile() is before
|
||||||
|
|
||||||
|
|
||||||
|
class TestDisplayManager:
|
||||||
|
@pytest.fixture
|
||||||
|
def dm(self):
|
||||||
|
from src.display_manager import DisplayManager
|
||||||
|
DisplayManager._instance = None
|
||||||
|
DisplayManager._initialized = False
|
||||||
|
manager = DisplayManager({"display": {
|
||||||
|
"hardware": {"rows": 32, "cols": 64, "chain_length": 1, "parallel": 1},
|
||||||
|
"runtime": {"gpio_slowdown": 0}}}, suppress_test_pattern=True)
|
||||||
|
yield manager
|
||||||
|
manager.render_gate = None
|
||||||
|
manager.set_scrolling_state(False)
|
||||||
|
DisplayManager._instance = None
|
||||||
|
DisplayManager._initialized = False
|
||||||
|
|
||||||
|
def test_opens_around_each_swap(self, dm):
|
||||||
|
calls = []
|
||||||
|
|
||||||
|
class Spy:
|
||||||
|
def before_swap(self, hold):
|
||||||
|
calls.append(("before", hold))
|
||||||
|
|
||||||
|
def after_swap(self, hold):
|
||||||
|
calls.append(("after", hold))
|
||||||
|
|
||||||
|
real_swap = dm.matrix.SwapOnVSync
|
||||||
|
|
||||||
|
def swap(canvas, *args, **kwargs):
|
||||||
|
calls.append(("swap",))
|
||||||
|
return real_swap(canvas, *args, **kwargs)
|
||||||
|
dm.matrix.SwapOnVSync = swap
|
||||||
|
dm.render_gate = Spy()
|
||||||
|
dm.set_scrolling_state(True, 2)
|
||||||
|
dm.draw.rectangle([0, 0, 3, 3], fill=(255, 0, 0))
|
||||||
|
dm.update_display()
|
||||||
|
assert calls == [("before", 2), ("swap",), ("after", 2)]
|
||||||
|
|
||||||
|
def test_off_screen_drawing_never_touches_it(self, dm):
|
||||||
|
class Boom:
|
||||||
|
def before_swap(self, hold):
|
||||||
|
raise AssertionError("an off-screen frame reached the gate")
|
||||||
|
after_swap = before_swap
|
||||||
|
|
||||||
|
dm.render_gate = Boom()
|
||||||
|
with dm.offscreen():
|
||||||
|
dm.draw.rectangle([0, 0, 3, 3], fill=(255, 0, 0))
|
||||||
|
dm.update_display()
|
||||||
|
|
||||||
|
|
||||||
|
class TestVegasWiring:
|
||||||
|
def _coordinator(self, **config):
|
||||||
|
from src.vegas_mode.config import VegasModeConfig
|
||||||
|
from src.vegas_mode.coordinator import VegasModeCoordinator
|
||||||
|
|
||||||
|
class Holder:
|
||||||
|
def __init__(self):
|
||||||
|
self._buffer_lock = threading.RLock()
|
||||||
|
self._prefetch_lock = threading.Lock()
|
||||||
|
self._cache_lock = threading.Lock()
|
||||||
|
|
||||||
|
c = VegasModeCoordinator.__new__(VegasModeCoordinator)
|
||||||
|
c.vegas_config = VegasModeConfig(**config)
|
||||||
|
c._state_lock = threading.Lock()
|
||||||
|
c.stream_manager = Holder()
|
||||||
|
c.render_pipeline = Holder()
|
||||||
|
c.plugin_adapter = Holder()
|
||||||
|
c.display_manager = type("DM", (), {"render_gate": None})()
|
||||||
|
return c
|
||||||
|
|
||||||
|
def test_off_by_default(self, monkeypatch):
|
||||||
|
monkeypatch.setattr(render_gate, "swap_releases_gil", lambda: True)
|
||||||
|
c = self._coordinator()
|
||||||
|
c._install_render_gate()
|
||||||
|
assert c.display_manager.render_gate is None
|
||||||
|
|
||||||
|
def test_installed_for_the_run_and_removed_after(self, monkeypatch):
|
||||||
|
monkeypatch.setattr(render_gate, "swap_releases_gil", lambda: True)
|
||||||
|
c = self._coordinator(prefetch_gate=True)
|
||||||
|
c._install_render_gate()
|
||||||
|
gate = c.display_manager.render_gate
|
||||||
|
assert isinstance(gate, RenderGate)
|
||||||
|
assert c._state_lock in gate._guarded
|
||||||
|
assert c.stream_manager._buffer_lock in gate._guarded
|
||||||
|
assert c.plugin_adapter._cache_lock in gate._guarded
|
||||||
|
c._remove_render_gate()
|
||||||
|
assert c.display_manager.render_gate is None
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("releases", [False, None])
|
||||||
|
def test_ignored_without_a_binding_that_releases_the_gil(self, monkeypatch, releases):
|
||||||
|
monkeypatch.setattr(render_gate, "swap_releases_gil", lambda: releases)
|
||||||
|
c = self._coordinator(prefetch_gate=True)
|
||||||
|
c._install_render_gate()
|
||||||
|
assert c.display_manager.render_gate is None
|
||||||
|
|
||||||
|
def test_read_from_config(self):
|
||||||
|
from src.vegas_mode.config import VegasModeConfig
|
||||||
|
on = VegasModeConfig.from_config(
|
||||||
|
{"display": {"vegas_scroll": {"prefetch_gate": True}}})
|
||||||
|
assert on.prefetch_gate is True
|
||||||
|
assert on.to_dict()["prefetch_gate"] is True
|
||||||
|
assert VegasModeConfig.from_config(
|
||||||
|
{"display": {"vegas_scroll": {}}}).prefetch_gate is False
|
||||||
|
|
||||||
|
def test_the_prefetch_runs_inside_the_gate(self):
|
||||||
|
from src.vegas_mode.render_pipeline import RenderPipeline
|
||||||
|
|
||||||
|
gate = RenderGate()
|
||||||
|
seen = []
|
||||||
|
|
||||||
|
class Stream:
|
||||||
|
def take_next_group(self, offscreen_only=False):
|
||||||
|
seen.append((sys.getprofile() == gate._hook, offscreen_only))
|
||||||
|
return ["segment"]
|
||||||
|
|
||||||
|
pipeline = RenderPipeline.__new__(RenderPipeline)
|
||||||
|
pipeline.config = type("C", (), {"continuous_scroll": True})()
|
||||||
|
pipeline._prefetch_lock = threading.Lock()
|
||||||
|
pipeline._prefetch_thread = None
|
||||||
|
pipeline._prepared_group = None
|
||||||
|
pipeline.stream_manager = Stream()
|
||||||
|
pipeline.display_manager = type("DM", (), {"render_gate": gate})()
|
||||||
|
pipeline.start_prefetch()
|
||||||
|
pipeline._prefetch_thread.join(5)
|
||||||
|
assert seen == [(True, True)]
|
||||||
|
assert pipeline._prepared_group == ["segment"]
|
||||||
Reference in New Issue
Block a user