diff --git a/src/background_data_service.py b/src/background_data_service.py index a41aef45..bd70d884 100644 --- a/src/background_data_service.py +++ b/src/background_data_service.py @@ -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. diff --git a/src/common/espn_dates.py b/src/common/espn_dates.py index bda7e86a..4a62f2cc 100644 --- a/src/common/espn_dates.py +++ b/src/common/espn_dates.py @@ -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,20 +199,25 @@ 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. """ - try: - response = session.get( - url, - params=dict(params, dates=chunk, limit=ESPN_MAX_LIMIT), - headers=headers, - timeout=timeout, - ) - response.raise_for_status() - return response_json(response) - except Exception as exc: # noqa: BLE001 - see docstring - if logger: - logger.warning("ESPN chunk %s failed, skipping it: %s", chunk, exc) - return None + with _yielding(): + try: + response = session.get( + url, + params=dict(params, dates=chunk, limit=ESPN_MAX_LIMIT), + headers=headers, + timeout=timeout, + ) + response.raise_for_status() + return response_json(response) + except Exception as exc: # noqa: BLE001 - see docstring + if logger: + logger.warning("ESPN chunk %s failed, skipping it: %s", chunk, exc) + return None def _fetch_chunks( diff --git a/src/common/render_gate.py b/src/common/render_gate.py index e1c81df3..a096db66 100644 --- a/src/common/render_gate.py +++ b/src/common/render_gate.py @@ -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() diff --git a/src/vegas_mode/coordinator.py b/src/vegas_mode/coordinator.py index 6713af02..009ed5d5 100644 --- a/src/vegas_mode/coordinator.py +++ b/src/vegas_mode/coordinator.py @@ -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) diff --git a/test/test_render_gate.py b/test/test_render_gate.py index a4481777..d7c06311 100644 --- a/test/test_render_gate.py +++ b/test/test_render_gate.py @@ -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):