mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-10 09:06:36 +00:00
feat(vegas): ESPN fetches give way to the render thread too
The hourly sports refresh froze hdpi's Vegas scroll for 1.9s: about twenty espn-chunk threads fetching and parsing at once, and the render thread (and the stall watchdog) queued behind all of them for the GIL. The prefetch gate only covered the prefetch thread. Vegas now makes its gate the active one (render_gate.set_active), and render_gate.yielding() gives way through it when there is one and does nothing otherwise. espn_dates wraps each chunk fetch in it (behind the same import fallback as json_body, for the copies plugins bundle), and the background data service wraps each worker. The render thread is never gated -- the first thread to swap is exempt, and a plugin pushing a live refresh from its update thread cannot take its place -- and nested blocks keep the outermost frame as the boundary for the lock checks. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -26,6 +26,7 @@ from enum import Enum
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
import pytz
|
||||
from src.cache_manager import CacheManager
|
||||
from src.common import render_gate
|
||||
from src.common.json_body import response_json
|
||||
from src.common.espn_dates import (
|
||||
RANGE_RETRY_SECONDS,
|
||||
@@ -317,6 +318,12 @@ class BackgroundDataService:
|
||||
return request_id
|
||||
|
||||
def _fetch_data_worker(self, request: FetchRequest) -> FetchResult:
|
||||
"""Fetch one request on a worker thread, giving way to Vegas's render
|
||||
thread while it scrolls (src/common/render_gate.py)."""
|
||||
with render_gate.yielding():
|
||||
return self._fetch_data(request)
|
||||
|
||||
def _fetch_data(self, request: FetchRequest) -> FetchResult:
|
||||
"""
|
||||
Worker function that performs the actual data fetching.
|
||||
|
||||
|
||||
@@ -47,6 +47,12 @@ except ImportError:
|
||||
def response_json(response: Any) -> Any:
|
||||
return response.json()
|
||||
|
||||
try:
|
||||
from src.common.render_gate import yielding as _yielding
|
||||
except ImportError:
|
||||
# Older cores (see above) have no render gate: fetch freely.
|
||||
from contextlib import nullcontext as _yielding
|
||||
|
||||
# Above this, ESPN returns a truncated list instead of an error. See module
|
||||
# docstring: 500 is the largest value measured to return complete data.
|
||||
ESPN_MAX_LIMIT = 500
|
||||
@@ -193,7 +199,12 @@ def _fetch_one_chunk(
|
||||
|
||||
One bad chunk must not sink the rest of the season, so every error is
|
||||
logged and swallowed here rather than raised to the gather below.
|
||||
|
||||
While Vegas scrolls, a chunk runs only when the render thread is waiting on
|
||||
the panel (``render_gate``): a season refresh is a couple of dozen of these
|
||||
at once, and they used to crowd the render thread off the GIL for seconds.
|
||||
"""
|
||||
with _yielding():
|
||||
try:
|
||||
response = session.get(
|
||||
url,
|
||||
|
||||
@@ -28,7 +28,16 @@ behind it, so it is never parked:
|
||||
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.
|
||||
whatever the gate gets wrong costs a frame, not a freeze. The render thread
|
||||
itself is never gated, whatever it calls.
|
||||
|
||||
The prefetch thread is not the only one that competes. Once an hour the sports
|
||||
plugins refresh their schedules together, about twenty ESPN chunk-fetch threads
|
||||
at once (hdpi, 2026-09-24), and the render thread queued behind all of them for
|
||||
a 1.9s freeze. So Vegas also makes its gate the *active* one, and code that runs
|
||||
background fetches wraps them in the module-level ``yielding()``, which uses the
|
||||
active gate if there is one and does nothing otherwise: ``espn_dates``' chunk
|
||||
fetches and the background data service's workers.
|
||||
|
||||
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
|
||||
@@ -42,7 +51,8 @@ import sys
|
||||
import threading
|
||||
import time
|
||||
from collections import deque
|
||||
from typing import Any, Callable, Deque, List, Optional
|
||||
from contextlib import nullcontext
|
||||
from typing import Any, Callable, ContextManager, 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.
|
||||
@@ -122,6 +132,7 @@ class RenderGate:
|
||||
self._period: Optional[float] = None
|
||||
self._guarded: List[Any] = []
|
||||
self._local = threading.local()
|
||||
self._render_ident: Optional[int] = None
|
||||
#: How often, and for how long in all, background threads were parked.
|
||||
self.parks = 0
|
||||
self.parked_seconds = 0.0
|
||||
@@ -164,6 +175,11 @@ class RenderGate:
|
||||
"""The swap returned and the render thread needs the GIL: close the window."""
|
||||
now = self.clock()
|
||||
self._open_until = 0.0
|
||||
if self._render_ident is None:
|
||||
# The first thread to swap is the render loop. A plugin pushing a
|
||||
# live refresh from its update thread swaps too, but must not take
|
||||
# over its exemption.
|
||||
self._render_ident = threading.get_ident()
|
||||
last = self._last_return
|
||||
if last is not None and now - last < STALE_SECONDS:
|
||||
self._periods.append((now - last) / max(1, int(hold)))
|
||||
@@ -208,15 +224,49 @@ class _Yielding:
|
||||
self.gate = gate
|
||||
self._previous: Any = None
|
||||
self._previous_base: Any = None
|
||||
self._skipped = False
|
||||
|
||||
def __enter__(self) -> RenderGate:
|
||||
local = self.gate._local # pylint: disable=protected-access
|
||||
gate = self.gate
|
||||
# pylint: disable=protected-access
|
||||
if threading.get_ident() == gate._render_ident:
|
||||
self._skipped = True # parking the render thread parks the display
|
||||
return gate
|
||||
local = gate._local
|
||||
self._previous_base = getattr(local, "base", None)
|
||||
local.base = sys._getframe(1) # pylint: disable=protected-access
|
||||
if self._previous_base is None:
|
||||
# Nested blocks keep the outermost frame, so everything the thread
|
||||
# entered since it first gave way is still checked for locks.
|
||||
local.base = sys._getframe(1)
|
||||
self._previous = sys.getprofile()
|
||||
sys.setprofile(self.gate._hook) # pylint: disable=protected-access
|
||||
return self.gate
|
||||
sys.setprofile(gate._hook)
|
||||
return gate
|
||||
|
||||
def __exit__(self, *_exc: Any) -> None:
|
||||
if self._skipped:
|
||||
return
|
||||
sys.setprofile(self._previous)
|
||||
self.gate._local.base = self._previous_base # pylint: disable=protected-access
|
||||
|
||||
|
||||
#: The gate of the Vegas run in progress, if any; see the module docstring.
|
||||
_active: Optional[RenderGate] = None
|
||||
|
||||
|
||||
def set_active(gate: Optional[RenderGate]) -> None:
|
||||
"""Make ``gate`` the one ``yielding()`` uses (None: no gate, run freely)."""
|
||||
global _active # pylint: disable=global-statement
|
||||
_active = gate
|
||||
|
||||
|
||||
def active() -> Optional[RenderGate]:
|
||||
"""The gate ``yielding()`` currently uses, if any."""
|
||||
return _active
|
||||
|
||||
|
||||
def yielding() -> ContextManager[Any]:
|
||||
"""``with render_gate.yielding():`` gives way to the render thread while a
|
||||
Vegas run has a gate, and does nothing otherwise. For background work that
|
||||
does not know whether Vegas is running."""
|
||||
gate = _active
|
||||
return gate.yielding() if gate is not None else nullcontext()
|
||||
|
||||
@@ -357,6 +357,7 @@ class VegasModeCoordinator:
|
||||
getattr(self.render_pipeline, '_prefetch_lock', None),
|
||||
getattr(self.plugin_adapter, '_cache_lock', None))
|
||||
self.display_manager.render_gate = gate
|
||||
render_gate.set_active(gate)
|
||||
logger.info("Vegas: prefetch gated on vsync")
|
||||
|
||||
def _remove_render_gate(self) -> None:
|
||||
@@ -364,6 +365,8 @@ class VegasModeCoordinator:
|
||||
if gate is None:
|
||||
return
|
||||
self.display_manager.render_gate = None
|
||||
if render_gate.active() is gate:
|
||||
render_gate.set_active(None)
|
||||
logger.info("Vegas: prefetch gate parked the prefetch %d times, %.1fs in all",
|
||||
gate.parks, gate.parked_seconds)
|
||||
|
||||
|
||||
@@ -242,6 +242,102 @@ class TestYielding:
|
||||
assert sys.getprofile() is before
|
||||
|
||||
|
||||
class TestWhoIsGated:
|
||||
def _in_thread(self, fn):
|
||||
out = []
|
||||
thread = threading.Thread(target=lambda: out.append(fn()))
|
||||
thread.start()
|
||||
thread.join(5)
|
||||
return out[0]
|
||||
|
||||
def test_the_render_thread_never_gives_way(self):
|
||||
gate = RenderGate()
|
||||
gate.before_swap(1)
|
||||
gate.after_swap(1) # this thread swaps: it is the render thread
|
||||
with gate.yielding():
|
||||
assert sys.getprofile() is not gate._hook
|
||||
|
||||
def background():
|
||||
with gate.yielding():
|
||||
return sys.getprofile() == gate._hook
|
||||
assert self._in_thread(background)
|
||||
|
||||
def test_a_live_refresh_from_another_thread_does_not_take_its_place(self):
|
||||
gate = RenderGate()
|
||||
gate.after_swap(1) # the render loop
|
||||
self._in_thread(lambda: gate.after_swap(1)) # a plugin pushing a frame
|
||||
with gate.yielding():
|
||||
assert sys.getprofile() is not gate._hook
|
||||
|
||||
def test_nested_blocks_keep_the_outer_boundary(self):
|
||||
gate = RenderGate()
|
||||
|
||||
def nested():
|
||||
with gate.yielding():
|
||||
outer = gate._local.base
|
||||
with gate.yielding():
|
||||
inner = gate._local.base
|
||||
after = gate._local.base
|
||||
still_hooked = sys.getprofile() == gate._hook
|
||||
return outer is inner is after and still_hooked, sys.getprofile()
|
||||
kept, final = self._in_thread(nested)
|
||||
assert kept and final is None
|
||||
|
||||
|
||||
class TestActiveGate:
|
||||
@pytest.fixture(autouse=True)
|
||||
def no_gate_left_behind(self):
|
||||
yield
|
||||
render_gate.set_active(None)
|
||||
|
||||
def test_without_a_vegas_run_nothing_is_gated(self):
|
||||
render_gate.set_active(None)
|
||||
with render_gate.yielding():
|
||||
assert sys.getprofile() is None
|
||||
|
||||
def test_with_one_background_work_gives_way(self):
|
||||
gate = RenderGate()
|
||||
render_gate.set_active(gate)
|
||||
with render_gate.yielding():
|
||||
assert sys.getprofile() == gate._hook
|
||||
assert sys.getprofile() is None
|
||||
|
||||
def test_espn_chunk_fetches_give_way(self):
|
||||
from src.common import espn_dates
|
||||
gate = RenderGate()
|
||||
render_gate.set_active(gate)
|
||||
seen = []
|
||||
|
||||
class Response:
|
||||
content = b'{"events": []}'
|
||||
|
||||
def raise_for_status(self):
|
||||
pass
|
||||
|
||||
def json(self):
|
||||
return {"events": []}
|
||||
|
||||
class Session:
|
||||
def get(self, *args, **kwargs):
|
||||
seen.append(sys.getprofile() == gate._hook)
|
||||
return Response()
|
||||
assert espn_dates._fetch_one_chunk(
|
||||
Session(), "http://x", {}, {}, 5, None, "20260924") == {"events": []}
|
||||
assert seen == [True]
|
||||
|
||||
def test_background_data_fetches_give_way(self):
|
||||
from src.background_data_service import BackgroundDataService
|
||||
gate = RenderGate()
|
||||
render_gate.set_active(gate)
|
||||
service = BackgroundDataService.__new__(BackgroundDataService)
|
||||
service._shutdown = True # never started: nothing for __del__ to stop
|
||||
seen = []
|
||||
service._fetch_data = lambda request: seen.append(
|
||||
sys.getprofile() == gate._hook) or "result"
|
||||
assert service._fetch_data_worker(object()) == "result"
|
||||
assert seen == [True]
|
||||
|
||||
|
||||
class TestDisplayManager:
|
||||
@pytest.fixture
|
||||
def dm(self):
|
||||
@@ -332,8 +428,10 @@ class TestVegasWiring:
|
||||
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
|
||||
assert render_gate.active() is gate # background fetches use it too
|
||||
c._remove_render_gate()
|
||||
assert c.display_manager.render_gate is None
|
||||
assert render_gate.active() is None
|
||||
|
||||
@pytest.mark.parametrize("releases", [False, None])
|
||||
def test_ignored_without_a_binding_that_releases_the_gil(self, monkeypatch, releases):
|
||||
|
||||
Reference in New Issue
Block a user