From bb23a29e8e989db264714921ab883003a228129f Mon Sep 17 00:00:00 2001 From: Chadwick Boulay Date: Tue, 8 Sep 2026 22:51:49 -0400 Subject: [PATCH 1/2] Rename chunk_dim to stream_dim, following ezmsg core Tracks the upstream rename of `AxisArray.chunk_dim`. `chunk` was overloaded -- xarray's `DataArray.chunk`, dask, zarr and HDF5 all use it for memory/IO layout, and ezmsg's own `array_chunker` / `chunk_len` mean "how much data per message" -- so the field now names its role instead: the dimension the stream extends along. Bumps the ezmsg floor to 3.10.0b3, the first release carrying the new name. Pure rename otherwise. --- pyproject.toml | 4 ++-- src/ezmsg/xdf/iter.py | 2 +- tests/test_iter.py | 18 +++++++++--------- 3 files changed, 12 insertions(+), 12 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 45243a7..6d8d6fb 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -8,8 +8,8 @@ authors = [ requires-python = ">=3.10.15" dynamic = ["version"] dependencies = [ - # 3.10.0b2 for AxisArray.chunk_dim and CoordinateAxis.fingerprint. - "ezmsg>=3.10.0b2", + # 3.10.0b3 for AxisArray.stream_dim and CoordinateAxis.fingerprint. + "ezmsg>=3.10.0b3", "ezmsg-baseproc>=1.11.0", "numpy>=2.0.2", "pyxdf>=1.16.8", diff --git a/src/ezmsg/xdf/iter.py b/src/ezmsg/xdf/iter.py index e204351..4379cf0 100644 --- a/src/ezmsg/xdf/iter.py +++ b/src/ezmsg/xdf/iter.py @@ -229,7 +229,7 @@ def _build_template(stream: dict, name: str, n_ch: int, fs: float) -> AxisArray: key=name, # Messages accumulate along `time`, whether the stream is regular or # carries per-sample timestamps; `ch` describes the stream itself. - chunk_dim="time", + stream_dim="time", ) diff --git a/tests/test_iter.py b/tests/test_iter.py index deb67c3..8a7ad17 100644 --- a/tests/test_iter.py +++ b/tests/test_iter.py @@ -109,7 +109,7 @@ def test_force_single_sample_splits_an_irregular_stream(self, test_xdf_path): class TestMessagesArriveReadyForConsumers: """Two things only the source can supply, both set once per stream. - ``chunk_dim`` names the dimension messages accumulate along -- the one whose + ``stream_dim`` names the dimension messages accumulate along -- the one whose length is just however much of the file this chunk covered, and which a consumer must leave out of the state it caches against the stream's configuration. ``fingerprint`` is the channel axis's content digest, cached @@ -117,13 +117,13 @@ class TestMessagesArriveReadyForConsumers: consumer in every process from recomputing it on every message. """ - def test_the_single_stream_iterator_declares_its_chunk_dim(self, test_xdf_path): - assert all(m.chunk_dim == "time" for m in eeg_messages(test_xdf_path)) + def test_the_single_stream_iterator_declares_its_stream_dim(self, test_xdf_path): + assert all(m.stream_dim == "time" for m in eeg_messages(test_xdf_path)) def test_the_multi_stream_iterator_declares_it_for_every_stream(self, test_xdf_path): it = XDFMultiAxArrIterator(filepath=test_xdf_path, chunk_dur=1.0) - undeclared = sorted({m.key for m in it if m is not None and m.chunk_dim != "time"}) - assert not undeclared, f"streams not declaring chunk_dim='time': {undeclared}" + undeclared = sorted({m.key for m in it if m is not None and m.stream_dim != "time"}) + assert not undeclared, f"streams not declaring stream_dim='time': {undeclared}" def test_the_channel_axis_is_primed(self, test_xdf_path): msg = eeg_messages(test_xdf_path)[0] @@ -143,7 +143,7 @@ def test_one_axis_object_serves_the_whole_stream(self, test_xdf_path): def test_the_chunk_axis_is_left_cold(self, test_xdf_path): """Digesting per-message timestamps would be pure cost: no consumer reads - the chunk axis's fingerprint.""" + the stream axis's fingerprint.""" it = XDFMultiAxArrIterator(filepath=test_xdf_path, chunk_dur=1.0) markers = [m for m in it if m is not None and m.key == MARKER_STREAM.name] assert markers, "no marker messages" @@ -152,7 +152,7 @@ def test_the_chunk_axis_is_left_cold(self, test_xdf_path): def test_it_all_survives_the_transport(self, test_xdf_path): msg = eeg_messages(test_xdf_path)[0] landed = pickle.loads(pickle.dumps(msg)) - assert landed.chunk_dim == "time" + assert landed.stream_dim == "time" assert "_fingerprint" in landed.axes["ch"].__dict__ assert landed.axes["ch"].__dict__["_fingerprint"] == msg.axes["ch"].fingerprint @@ -292,7 +292,7 @@ def test_the_single_stream_unit_publishes_the_whole_file(self, test_xdf_path): ) assert msgs, "no messages published" assert sum(m.data.shape[0] for m in msgs) == EEG_STREAM.n_samples - assert all(m.chunk_dim == "time" for m in msgs) + assert all(m.stream_dim == "time" for m in msgs) assert all("_fingerprint" in m.axes["ch"].__dict__ for m in msgs) def test_the_multi_stream_unit_publishes_both_streams(self, test_xdf_path): @@ -303,4 +303,4 @@ def test_the_multi_stream_unit_publishes_both_streams(self, test_xdf_path): assert {m.key for m in msgs} == {EEG_STREAM.name, MARKER_STREAM.name} eeg = [m for m in msgs if m.key == EEG_STREAM.name] assert sum(m.data.shape[0] for m in eeg) == EEG_STREAM.n_samples - assert all(m.chunk_dim == "time" for m in msgs) + assert all(m.stream_dim == "time" for m in msgs) From 239680a3521b380439180fd541801cfde61c1f6f Mon Sep 17 00:00:00 2001 From: Chadwick Boulay Date: Wed, 9 Sep 2026 03:01:07 -0400 Subject: [PATCH 2/2] Cap the test job at 30 minutes and supersede in-flight runs A wedged job used to run to GitHub's 6-hour default. One did: an ezmsg-simbiophys run hung immediately after `collected 161 items`, produced not one test dot, and sat on a runner for 3h09m until it was cancelled by hand -- long enough to starve the macOS queue for every other repo in the org. `timeout-minutes: 30` bounds that; the suite runs in well under a minute on every platform, so the margin is generous. The concurrency group makes a new push to a PR cancel the run already in flight for that ref, so a stuck job cannot hold a slot while its own replacement queues behind it. --- .github/workflows/python-tests.yml | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/.github/workflows/python-tests.yml b/.github/workflows/python-tests.yml index 57abf79..1c6ac98 100644 --- a/.github/workflows/python-tests.yml +++ b/.github/workflows/python-tests.yml @@ -9,6 +9,12 @@ on: - dev workflow_dispatch: +# A new push to a PR supersedes the run already in flight for that ref, so a +# wedged job cannot sit on a runner while its replacement queues behind it. +concurrency: + group: ${{ github.workflow }}-${{ github.ref }} + cancel-in-progress: true + jobs: build: strategy: @@ -19,6 +25,7 @@ jobs: - "windows-latest" - "macos-latest" runs-on: ${{matrix.os}} + timeout-minutes: 30 steps: - uses: actions/checkout@v4