Files
LEDMatrix/test/test_vegas_live_worker.py
ChuckandClaude Opus 5.5 795834811f perf(scroll): extend and trim the Vegas strip in place (#701)
* perf(timing): say which render-thread work a late frame followed

The soak already says how often a moving frame reached the panel late, but
not what the render thread was doing just before it. Vegas does two kinds of
work there between frames -- building its strip (compose, extend) and, with
live elements, patching changed pixels into it -- and deciding whether either
is affordable needs their own numbers.

- FrameTimingRecorder.note_op(kind, nbytes) tags the next presented frame.
  Totals gain op_frames, late_op_frames, op_freezes and op_bytes per kind;
  aggregate() still takes frames without ops. The file schema is unchanged.
- Vegas tags compose and every strip extension (with the bytes it copied).
- frame_soak prints an "after work" table: frames, late %, freezes and MB
  moved per kind, only when something tagged its work.
- render_bench gains --strip-screens (Vegas-sized strips), --patch-bytes /
  --patch-every / --patch-where (in-place column writes, as a live element
  update does) and --extend-every-screens / --extend-width (append + trim on
  a fixed cadence that holds the strip's width).

No runtime behaviour changes: this is the measurement gate for live Vegas
elements.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* docs(changelog): note the frame-op attribution and bench modes

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* perf(scroll): build the strip's PIL image only when something reads it

Every Vegas strip extension rebuilt ScrollHelper.cached_image from
cached_array in full, twice (append, then trim), on the render thread:
Image.fromarray is 1.7ms for an 8,000px strip and 3.8ms for 20,000px on a
Pi 4 (measured on ledpi), about two thirds of an extension's render-thread
cost. Nothing on the frame path reads the image's pixels; every frame is cut
from the array.

cached_image is now a property. append_content and drop_scrolled_prefix
defer it; the first read builds it from the array it started with and keeps
it only if the strip has not changed meanwhile, so a sync push racing an
extension cannot leave a stale image cached. Assigning cached_image stores
exactly what was assigned, as before. has_strip() says whether there is a
strip without building its image; the helper's frame path, Vegas and the
adapter's scroll-cache invalidation use it. The strip is also no longer held
in memory twice.

In Vegas the image is now built only by a multi-display sync push.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* feat(vegas): live elements -- a plugin API for content that changes while it scrolls

Vegas bakes each plugin's pictures into one strip, so a card already on its
way across the panel keeps what it showed when it was drawn. This adds the
API and bookkeeping for content that can be updated in place; the worker
that redraws and swaps it follows separately. No shipped plugin implements
the hook yet, so nothing changes for users.

Plugin API (core 3.8.0), all no-ops by default:
- BasePlugin.get_vegas_elements() -> [VegasElement(key, image, version,
  live, refresh_hz)]: named, fixed-width pieces of Vegas content.
- BasePlugin.redraw_vegas_element(key, width, height, at): a lock-free
  redraw for content that changes with time.
- BasePlugin.notify_vegas_data_changed(): data that lands outside update().
- src/plugin_system/vegas_elements.py (VegasElement, re-exported from
  base_plugin).

Core:
- PluginAdapter asks a plugin that implements the hook for elements on the
  background fetch only (under its lock, on its own canvas); every other
  path keeps get_vegas_content(). Live elements are pinned (padded with
  content_padding, never trimmed), tagged with their key, digest and data
  epoch in Image.info so the existing cache and group plumbing carry them
  unchanged, and untagged if a width budget crops them.
- RenderPipeline records where each live element lands (ElementRecord), in
  absolute strip columns a trim does not move; the block-start arithmetic
  is shared with the STATIC markers.
- PluginManager update listeners (add/remove_update_listener,
  notify_data_changed): told the moment update() completes, not at the
  next ~4s Vegas poll. The coordinator uses one to move each plugin's data
  epoch on.
- vegas_scroll.live_refresh (kill switch), live_max_hz, live_min_interval,
  live_lead_screens; per-plugin core-owned vegas_live. Live elements are
  off under multi-display sync, in swap mode and with offscreen_prefetch off.
- scripts/check_plugin.py checks the element contract
  (src/plugin_system/testing/vegas.py); test/fixtures/plugins/vegas-live-stub
  is a working example.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* 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>

* feat(sports): live Vegas cards for the scoreboards (shared layer)

One live element per game, drawn only when what the card shows changes, so
a score changes on a card already crossing the panel. The shared part, so
each scoreboard adopts it in a few lines:

- src/common/sports_vegas.py: game_key, game_fingerprint (the whole game
  dict, frozen: no drawn field can be missed), dedupe_games, VegasCardCache,
  StickyOdds (odds a live poll left out stay drawn), finished_games /
  with_finished_games (a game that just went final keeps its card, after its
  league's live games; one a heuristic only judged over keeps its live
  state, so a tied end of regulation never shows FINAL early).
- SportsScrollDisplay.make_vegas_renderer() is the override point;
  build_vegas_elements() and SportsScrollDisplayManager
  .get_vegas_elements_for() do the rest. A card's version includes its
  teams' ranks, which the renderer draws from the rankings cache.
- SportsLiveSharedMixin._record_finished_game() / finished_games_snapshot():
  held for FINISHED_GAME_TTL after it leaves the live list.

A sport that does not implement make_vegas_renderer keeps its ordinary Vegas
content, so no scoreboard changes until it opts in.

scripts/render_plugin.py --vegas renders a plugin's Vegas block as the
ticker lays it out, and --timeline stacks it at successive moments as
the ticker would update it in place; the join is now
render_pipeline.join_plugin_rows().

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* feat(vegas): keep live games in the ticker by default

display.vegas_scroll.live_in_ticker now defaults to true: through a live
game the marquee keeps running and the live scoreboard takes extra turns in
it -- its cards updating in place while they scroll -- instead of the ticker
giving way to the full-screen scoreboard.

The new default would reach nobody on its own: every existing config holds
an explicit false copied from the template (there was no control for it),
and the template merge only adds missing keys. ConfigManager therefore turns
a stored false on once, with a backup, and records live_in_ticker_migrated
so a false chosen afterwards stays. The marker is never in the template.

A "Keep live games in the ticker" checkbox under Vegas mode sets it. Tests
that pin the full-screen takeover now say live_in_ticker=false.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* feat(dev): preview a plugin's Vegas strip in the dev server

The dev server's View selector gains "Vegas strip (live elements)" and
"Vegas strip (plain Vegas content)": the plugin's block of the Vegas ticker,
laid out by the ticker's own code (render_vegas_strip, as render_plugin.py
--vegas uses), with its live elements listed. /api/render takes
"vegas": "live" | "plain"; the display view is unchanged.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* perf(scroll): extend and trim the Vegas strip in place

Every strip extension rebuilt the whole strip (np.concatenate: 2-2.6 ms for
a 10-14k px strip at 512x64 on a Pi 4) and every trim copied what was left
(1.2-1.8 ms), on the render thread. With ~4 ms of slack per refresh, every
extension frame on hdpi missed its refresh (5/5 in each soak run).

The strip now lives in a buffer with spare room; cached_array is a view of
its live columns. An append writes only the new columns (~0.2 ms), a trim
only moves the view's start, and the one full copy happens when the buffer
is reallocated (STRIP_SPARE_FACTOR 3: about once every two strip-lengths
scrolled). A cached_array set from outside -- the multi-display follower's
read-only one, create_scrolling_image's -- is never written through, and a
new strip lets the old buffer go. last_copy_bytes says what was copied, and
the Vegas frame-timing attribution reports that instead of the whole strip.

test_scroll_helper_in_place.py: the buffer is reused and only new columns
copied, trims copy nothing, reallocation when the room runs out, outside
arrays untouched, and random appends/trims/patches/scrolling checked frame
by frame against the old copying strip (mutation-checked).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* refactor(sports): a default _determine_game_type on SportsScrollDisplay

render_vegas_card looked the method up with getattr and a None default, which
static analysis (Codacy) reports as calling something that may not be
callable. The base class now has the default -- the card type from the game's
state -- and the plugins that define their own override it as before.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* refactor(scroll): no assert in _extended_strip

An assert vanishes under python -O (Codacy); a real check says the same.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* fix: review follow-ups on the shared live-card layer

- The reused Vegas renderer always gets the current rankings, empty
  included, so ranks cleared since are not kept drawn.
- render_plugin.py: --timeline refuses --no-live (a timeline shows live
  elements changing), --timeline/--no-live need --vegas, and the Vegas
  paths create the output's directory like the display path does.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

* perf(vegas): lay a group's blocks out before the extension, off the render thread

#701 cut the strip copy, but the frame after an extension was still late:
the render thread also joined each plugin's rows (separation_gap measures
every pair) and pasted the blocks into an addition image, ~37 ms on hdpi
(Pi 4, 512x64) against ~3.75 ms of slack.

The thread that fetched the group now does that as each member arrives:
RenderPipeline.prepare_group_member joins the rows and turns the block into
pixels (the prefetch thread and the live worker, under the render gate).
extend_scroll_content takes those blocks, and ScrollHelper.append_content
writes items -- images or RGB arrays -- straight into the strip's spare room,
blanking only the gaps. A member that was not prepared (an inline fetch) is
joined at the extension as before; the strip is identical either way.

On hdpi the extension's render-thread work goes from 37.5 ms to 3.2 ms p50
in place (6.9 ms when the buffer is reallocated). The helper's per-append
INFO line, a duplicate of the pipeline's, is now DEBUG.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-01 09:48:47 -04:00

496 lines
17 KiB
Python

"""The live-element worker's choices and hand-overs (src/vegas_mode/live_worker.py).
The worker is driven here one decision at a time -- _pick() then _run() --
against a fake pipeline, so every rule can be pinned without threads or
timing: what runs first, what is skipped, what reaches the render thread.
"""
import collections
import sys
import threading
from pathlib import Path
from types import SimpleNamespace
import numpy as np
import pytest
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from src.vegas_mode import live_worker # noqa: E402
from src.vegas_mode.config import VegasModeConfig # noqa: E402
from src.vegas_mode.elements import ( # noqa: E402
ElementRecord, LiveEpochs, LiveView, RenderedElement, pixel_digest,
)
from src.vegas_mode.live_worker import VegasWorker # noqa: E402
W, H = 128, 32
NOW = 1000.0
def _pixels(width, value):
array = np.full((H, width, 3), value, dtype=np.uint8)
array.setflags(write=False)
return array
def _element(key, width, value, epoch=1):
pixels = _pixels(width, value)
return RenderedElement(key=key, epoch=epoch, version=value, pixels=pixels,
digest=pixel_digest(pixels), width=width)
def _record(seq, key, abs_x, width=40, pid="p", epoch=0, value=0, hz=0.0):
return ElementRecord(seq=seq, plugin_id=pid, key=key, abs_x=abs_x, width=width,
epoch=epoch, digest=pixel_digest(_pixels(width, value)),
refresh_hz=hz)
class _Adapter:
def __init__(self):
self.live_epochs = LiveEpochs()
self.batches = {} # pid -> {key: RenderedElement}
self.redraws = {} # key -> RenderedElement | None
self.busy = set()
self.calls = []
def render_live_elements(self, plugin, pid, lock_timeout):
self.calls.append(("render", pid, lock_timeout))
if pid in self.busy:
return None
return self.live_epochs.get(pid), self.batches.get(pid, {})
def redraw_live_element(self, plugin, pid, key, width, height, at):
self.calls.append(("redraw", pid, key, width, at))
return self.redraws.get(key)
has_redraw = True
def has_lock_free_redraw(self, plugin):
return self.has_redraw
class _Stream:
def __init__(self, adapter):
self.plugin_adapter = adapter
self.plugin_manager = SimpleNamespace(plugins={"p": object(), "q": object()})
self.plans = []
self.fetched = []
def plan_next_group(self, count=None):
self.plans.append(count)
return ["a", "b", "c"]
def fetch_group_member(self, pid, offscreen_only=False):
self.fetched.append((pid, offscreen_only))
return (pid, [f"img-{pid}"])
def _pipeline(gate=True, **cfg):
adapter = _Adapter()
p = SimpleNamespace(
_strip_gen=1, _view=None, _elements=(), _prepared_group=None,
_prefetch_lock=threading.Lock(), _prefetch_generation=7, _applied={},
_live_slots={}, _live_ready=collections.deque(),
display_width=W, display_height=H, frame_interval=0.01,
config=VegasModeConfig(**cfg),
display_manager=SimpleNamespace(render_gate=_Gate() if gate else None),
stream_manager=_Stream(adapter), _prefetch_thread=None, prepared=[])
p.prepare_group_member = p.prepared.append
return p, adapter
class _Gate:
def __init__(self):
self.entered = 0
def yielding(self):
gate = self
class _Ctx:
def __enter__(self):
gate.entered += 1
def __exit__(self, *exc):
return False
return _Ctx()
def _view(left=1000, end=5000, t=NOW):
return LiveView(abs_left=left, abs_right=left + W, abs_end=end, t_mono=t)
def _worker(p, records=(), view=None, clock=lambda: NOW):
p._elements = tuple(records)
for r in records:
p._applied[r.seq] = (r.epoch, r.digest)
p._view = view if view is not None else _view()
return VegasWorker(p, clock=clock)
# -- choosing -----------------------------------------------------------------
def test_an_urgent_group_goes_before_anything():
p, adapter = _pipeline()
visible = _record(1, "k", 1010)
worker = _worker(p, [visible], _view(left=1000, end=1000 + W + 100))
adapter.live_epochs.bump("p")
worker.request_group()
assert worker._pick(NOW) == ("group", None)
def test_visible_data_then_ticks_then_a_normal_group_then_data_ahead():
p, adapter = _pipeline()
visible = _record(1, "vis", 1010)
animated = _record(2, "map", 1060, hz=4)
ahead = _record(3, "far", 3000, pid="q")
worker = _worker(p, [visible, animated, ahead])
adapter.live_epochs.bump("p")
adapter.live_epochs.bump("q")
worker.request_group()
assert worker._pick(NOW) == ("data", "p")
worker._handled_epoch[1] = worker._handled_epoch[2] = adapter.live_epochs.get("p")
assert worker._pick(NOW) == ("tick", animated)
worker._next_tick[2] = NOW + 10
assert worker._pick(NOW) == ("group", None)
worker._group_wanted = False
assert worker._pick(NOW) == ("data", "q")
def test_nothing_behind_the_viewport_is_redrawn():
p, adapter = _pipeline()
behind = _record(1, "gone", 900, width=40)
worker = _worker(p, [behind])
adapter.live_epochs.bump("p")
assert worker._pick(NOW) is None
def test_only_group_work_while_frames_have_stopped():
p, adapter = _pipeline()
worker = _worker(p, [_record(1, "k", 1010, hz=4)], _view(t=NOW - 5))
adapter.live_epochs.bump("p")
assert worker._pick(NOW) is None
worker.request_group()
assert worker._pick(NOW) == ("group", None)
def test_an_element_placed_from_older_data_is_caught_up():
# A group drawn before an update() and placed after it: its records carry
# the old epoch, so they are due at once.
p, adapter = _pipeline()
adapter.live_epochs.bump("p")
adapter.live_epochs.bump("p")
worker = _worker(p, [_record(1, "k", 1010, epoch=1)])
assert worker._pick(NOW) == ("data", "p")
def test_the_data_floor_defers_but_does_not_drop():
p, adapter = _pipeline(live_min_interval=2.0)
worker = _worker(p, [_record(1, "k", 1010)])
adapter.live_epochs.bump("p")
worker._last_data_job["p"] = NOW - 1.0
assert worker._pick(NOW) is None
p._view = _view(t=NOW + 1.5)
assert worker._pick(NOW + 1.5) == ("data", "p")
# -- data refreshes ------------------------------------------------------------
def test_a_refresh_hands_over_only_what_changed():
p, adapter = _pipeline()
same = _record(1, "same", 1010, value=0)
changed = _record(2, "changed", 1060, value=0)
worker = _worker(p, [same, changed])
epoch = adapter.live_epochs.bump("p")
adapter.batches["p"] = {"same": _element("same", 40, 0, epoch),
"changed": _element("changed", 40, 9, epoch)}
worker._run(("data", "p"))
assert list(p._live_ready) == [2]
patch = p._live_slots[2]
assert patch.epoch == epoch and patch.strip_gen == 1
assert worker._handled_epoch == {1: epoch, 2: epoch}
assert worker._pick(NOW + 100) is None # nothing left due
assert adapter.calls[0] == ("render", "p", live_worker.DATA_LOCK_TIMEOUT)
def test_a_redraw_of_another_width_is_refused_and_said_once(caplog):
p, adapter = _pipeline()
worker = _worker(p, [_record(1, "k", 1010, width=40)])
for n in range(3):
epoch = adapter.live_epochs.bump("p")
adapter.batches["p"] = {"k": _element("k", 44, n + 1, epoch)}
with caplog.at_level("INFO"):
worker._run(("data", "p"))
assert not p._live_ready
assert worker.stats["refused"] == 3
assert sum("must not change" in r.message for r in caplog.records) == 1
def test_a_busy_lock_backs_off_instead_of_waiting():
p, adapter = _pipeline()
worker = _worker(p, [_record(1, "k", 1010)])
adapter.live_epochs.bump("p")
adapter.busy.add("p")
worker._run(("data", "p"))
assert worker.stats["lock_busy"] == 1
assert worker._backoff_until["p"] == NOW + live_worker.LOCK_BACKOFF_S
p._view = _view(t=NOW + 2.5)
assert worker._pick(NOW + 0.5) is None # backing off (and floored)
adapter.busy.clear()
assert worker._pick(NOW + 2.5) == ("data", "p") # tried again later
def test_a_key_the_plugin_dropped_keeps_its_pixels():
p, adapter = _pipeline()
worker = _worker(p, [_record(1, "gone", 1010)])
epoch = adapter.live_epochs.bump("p")
adapter.batches["p"] = {}
worker._run(("data", "p"))
assert not p._live_ready and worker._handled_epoch[1] == epoch
def test_the_latest_hand_over_wins_the_slot():
p, adapter = _pipeline()
worker = _worker(p, [_record(1, "k", 1010)])
for value in (5, 6):
epoch = adapter.live_epochs.bump("p")
adapter.batches["p"] = {"k": _element("k", 40, value, epoch)}
worker._last_data_job.clear()
worker._run(("data", "p"))
assert list(p._live_ready) == [1, 1]
assert p._live_slots[1].pixels[0, 0, 0] == 6
# -- ticks ----------------------------------------------------------------------
def test_a_tick_uses_the_lock_free_redraw_and_reschedules():
p, adapter = _pipeline(live_max_hz=5)
record = _record(1, "map", 1010, width=40, hz=4)
worker = _worker(p, [record])
adapter.redraws["map"] = _element("map", 40, 3)
worker._run(("tick", record))
assert adapter.calls[0][0] == "redraw"
assert list(p._live_ready) == [1]
assert worker._next_tick[1] > 0
def test_a_tick_without_a_redraw_never_waits_for_the_lock():
p, adapter = _pipeline()
record = _record(1, "map", 1010, hz=4)
worker = _worker(p, [record])
adapter.has_redraw = False
worker._run(("tick", record))
assert adapter.calls == [("render", "p", 0.0)]
def test_a_redraw_that_returns_none_is_a_skip_not_a_full_redraw():
# None is the plugin saying "nothing new"; answering it with a locked
# get_vegas_elements() at the tick rate is exactly what the hook avoids.
p, adapter = _pipeline()
record = _record(1, "map", 1010, hz=4)
worker = _worker(p, [record])
adapter.redraws["map"] = None
worker._run(("tick", record))
assert [c[0] for c in adapter.calls] == ["redraw"]
assert not p._live_ready
def test_animation_is_capped_without_the_gate_and_when_slow():
p, _ = _pipeline(gate=False, live_max_hz=5)
record = _record(1, "map", 1010, hz=4)
worker = _worker(p, [record])
assert worker._tick_hz(record) == live_worker.UNGATED_MAX_HZ
p2, _ = _pipeline(live_max_hz=5)
worker2 = _worker(p2, [record])
assert worker2._tick_hz(record) == 4
worker2._render_ewma[("p", "map")] = live_worker.SLOW_RENDER_S * 2
assert worker2._tick_hz(record) == 2
p3, _ = _pipeline(live_max_hz=0)
assert _worker(p3, [record])._tick_hz(record) == 0
def test_throttling_never_raises_a_slow_elements_rate():
p, _ = _pipeline(live_max_hz=5)
slow = _record(1, "map", 1010, hz=live_worker.MIN_THROTTLED_HZ / 2)
worker = _worker(p, [slow])
worker._render_ewma[("p", "map")] = live_worker.SLOW_RENDER_S * 2
assert worker._tick_hz(slow) == live_worker.MIN_THROTTLED_HZ / 2
def test_a_tick_left_by_an_element_trimmed_away_does_not_spin_the_worker():
p, _ = _pipeline()
worker = _worker(p, [_record(1, "map", 1010, hz=4)])
worker._next_tick[1] = NOW - 5 # past due
p._elements = () # ...and trimmed off the strip
assert worker._next_wait(NOW) == live_worker.IDLE_WAIT_S
def test_a_tick_for_an_element_no_longer_animated_is_dropped():
p, _ = _pipeline(live_max_hz=0)
record = _record(1, "map", 1010, hz=4)
worker = _worker(p, [record])
worker._next_tick[1] = NOW - 5
assert worker._next_wait(NOW) == live_worker.IDLE_WAIT_S
assert worker._pick(NOW) is None
assert 1 not in worker._next_tick
def test_the_wait_is_until_the_next_tick_on_screen():
p, _ = _pipeline(live_max_hz=5)
record = _record(1, "map", 1010, hz=4)
worker = _worker(p, [record])
worker._next_tick[1] = NOW + 0.2
assert worker._next_wait(NOW) == pytest.approx(0.2)
def test_ticks_stop_for_an_element_far_ahead():
p, _ = _pipeline(live_lead_screens=1.0)
far = _record(1, "map", 1000 + 3 * W, hz=4)
worker = _worker(p, [far])
worker._next_tick[1] = NOW
assert worker._pick(NOW) is None
assert 1 not in worker._next_tick
# -- groups ---------------------------------------------------------------------
def test_a_group_is_fetched_a_member_at_a_time_and_published():
p, _ = _pipeline()
worker = _worker(p)
worker.request_group()
worker.request_group() # coalesced
for _ in range(3):
assert p._prepared_group is None
worker._run(worker._pick(NOW))
assert p._prepared_group == [("a", ["img-a"]), ("b", ["img-b"]), ("c", ["img-c"])]
# Each member was laid out for the strip here, as it arrived, not by the
# render thread at the extension.
assert p.prepared == p._prepared_group
assert p.stream_manager.plans == [None]
assert all(offscreen for _pid, offscreen in p.stream_manager.fetched)
assert worker._pick(NOW) is None # the slot is full
def test_a_reset_mid_group_drops_it():
p, _ = _pipeline()
worker = _worker(p)
worker.request_group()
worker._run(worker._pick(NOW))
p._prefetch_generation += 1
worker._run(("group", None))
worker._run(("group", None))
assert p._prepared_group is None
def test_a_stopped_worker_hands_over_what_it_has_of_a_group():
p, _ = _pipeline()
worker = _worker(p)
worker.request_group()
worker._run(worker._pick(NOW)) # one member of three
worker.stop()
worker._hand_over_partial_group()
assert p._prepared_group == [("a", ["img-a"])]
@pytest.mark.parametrize("why", ["reset", "a group already waiting"])
def test_a_partial_group_is_not_handed_over(why):
p, _ = _pipeline()
worker = _worker(p)
worker.request_group()
worker._run(worker._pick(NOW))
if why == "reset":
p._prefetch_generation += 1
else:
p._prepared_group = ["waiting"]
worker._hand_over_partial_group()
assert p._prepared_group == (None if why == "reset" else ["waiting"])
def test_a_new_worker_waits_for_the_one_it_replaces():
import time
p, _ = _pipeline()
retired = threading.Thread(target=time.sleep, args=(0.05,))
retired.start()
p._retired_worker = retired
_worker(p)._join_legacy_prefetch()
assert not retired.is_alive()
def test_every_job_runs_inside_the_gate():
p, adapter = _pipeline()
worker = _worker(p, [_record(1, "k", 1010)])
adapter.live_epochs.bump("p")
worker._run(("data", "p"))
assert p.display_manager.render_gate.entered == 1
# -- robustness -------------------------------------------------------------------
def test_a_job_that_raises_is_counted_and_the_worker_goes_on():
p, adapter = _pipeline()
worker = _worker(p, [_record(1, "k", 1010)])
adapter.live_epochs.bump("p")
def explode(*a, **k):
raise KeyError("plugin bug")
adapter.render_live_elements = explode
worker._run(("data", "p"))
assert worker.stats["errors"] == 1
worker.request_group()
worker._run(worker._pick(NOW))
assert p.stream_manager.plans
def test_a_new_strip_forgets_the_old_ones_bookkeeping():
p, _ = _pipeline()
worker = _worker(p, [_record(1, "k", 1010)])
worker._handled_epoch[1] = 5
worker._next_tick[1] = NOW
p._strip_gen += 1
p._elements = ()
worker._pick(NOW)
assert worker._handled_epoch == {} and worker._next_tick == {}
def test_the_thread_starts_and_stops():
import time
p, _ = _pipeline()
worker = _worker(p, clock=time.monotonic)
worker.start()
worker.request_group()
for _ in range(200):
if p._prepared_group is not None:
break
time.sleep(0.01)
worker.stop()
worker.join(2)
assert not worker.is_alive()
assert p._prepared_group is not None
@pytest.mark.parametrize("view", [None])
def test_no_view_yet_still_fetches_groups(view):
p, _ = _pipeline()
worker = VegasWorker(p)
worker.request_group()
assert worker._pick(NOW) == ("group", None)
def test_the_summary_names_the_slowest_redraw(caplog):
clock = [NOW]
p, adapter = _pipeline()
record = _record(1, "map", 1010, hz=4)
worker = _worker(p, [record], clock=lambda: clock[0])
adapter.redraws["map"] = _element("map", 40, 3)
worker._run(("tick", record))
clock[0] += live_worker.SUMMARY_INTERVAL_S + 1
with caplog.at_level("INFO"):
worker._maybe_summarise()
line = next(r.getMessage() for r in caplog.records if "Vegas live:" in r.getMessage())
assert "slowest redraw" in line and "p 'map'" in line
assert worker._slowest_redraw is None # per summary interval