From b64b119da0b6801f211f035c745a4a5f6123db38 Mon Sep 17 00:00:00 2001 From: Chadwick Boulay Date: Tue, 8 Sep 2026 22:54:23 -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. 161 passed. --- pyproject.toml | 4 ++-- src/ezmsg/simbiophys/baseline_drift.py | 2 +- src/ezmsg/simbiophys/dnss/lfp.py | 2 +- src/ezmsg/simbiophys/dnss/spike.py | 2 +- src/ezmsg/simbiophys/dynamic_colored_noise.py | 2 +- src/ezmsg/simbiophys/line_noise.py | 2 +- src/ezmsg/simbiophys/noise.py | 4 ++-- src/ezmsg/simbiophys/oscillator.py | 4 ++-- tests/unit/test_state_reset_semantics.py | 12 ++++++------ 9 files changed, 17 insertions(+), 17 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 890bb1c..4d99710 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -10,10 +10,10 @@ requires-python = ">=3.10" dynamic = ["version"] dependencies = [ "array-api-compat>=1.11.0", - # 3.10.0b2 for AxisArray.chunk_dim, which every producer here declares, and + # 3.10.0b3 for AxisArray.stream_dim, which every producer here declares, and # CoordinateAxis.fingerprint, which they prime so downstream consumers do # not each recompute it. - "ezmsg>=3.10.0b2", + "ezmsg>=3.10.0b3", "ezmsg-baseproc>=1.12.0", # axis-aware default state hash + hash witness "ezmsg-sigproc>=2.23.0", "ezmsg-event>=0.9.0", diff --git a/src/ezmsg/simbiophys/baseline_drift.py b/src/ezmsg/simbiophys/baseline_drift.py index 80a79cb..573ca4d 100644 --- a/src/ezmsg/simbiophys/baseline_drift.py +++ b/src/ezmsg/simbiophys/baseline_drift.py @@ -167,7 +167,7 @@ class BaselineDriftTransformer( """ # No `_hash_message`: the default already folds in the channel count and the - # chunk axis's gain, which is all this hashed before, and additionally the + # stream axis's gain, which is all this hashed before, and additionally the # channel *fingerprint* -- one drift process is warmed up per channel, so a # relabel at a fixed count leaves each channel wearing another's drift. diff --git a/src/ezmsg/simbiophys/dnss/lfp.py b/src/ezmsg/simbiophys/dnss/lfp.py index cd30025..2914fc0 100644 --- a/src/ezmsg/simbiophys/dnss/lfp.py +++ b/src/ezmsg/simbiophys/dnss/lfp.py @@ -247,7 +247,7 @@ def _reset_state(self, time_axis: LinearAxis) -> None: data=np.zeros((0, self.settings.n_ch), dtype=np.float64), dims=["time", "ch"], axes={"time": time_axis, "ch": ch_axis}, - chunk_dim="time", + stream_dim="time", ) def _produce(self, n_samples: int, time_axis: LinearAxis) -> AxisArray: diff --git a/src/ezmsg/simbiophys/dnss/spike.py b/src/ezmsg/simbiophys/dnss/spike.py index 9a884e0..2f39d98 100644 --- a/src/ezmsg/simbiophys/dnss/spike.py +++ b/src/ezmsg/simbiophys/dnss/spike.py @@ -316,7 +316,7 @@ def _reset_state(self, time_axis: LinearAxis) -> None: ), dims=["time", "ch"], axes={"time": time_axis, "ch": ch_axis}, - chunk_dim="time", + stream_dim="time", ) def _produce(self, n_samples: int, time_axis: LinearAxis) -> AxisArray: diff --git a/src/ezmsg/simbiophys/dynamic_colored_noise.py b/src/ezmsg/simbiophys/dynamic_colored_noise.py index 0fde279..5e779c0 100644 --- a/src/ezmsg/simbiophys/dynamic_colored_noise.py +++ b/src/ezmsg/simbiophys/dynamic_colored_noise.py @@ -220,7 +220,7 @@ class DynamicColoredNoiseTransformer( """ # No `_hash_message`: the default already folds in the channel count and the - # chunk axis's gain, which is all this hashed before, and additionally the + # stream axis's gain, which is all this hashed before, and additionally the # channel *fingerprint* -- the delay lines and per-channel coefficients are # tied to specific channels, so a relabel at a fixed count has to reset. diff --git a/src/ezmsg/simbiophys/line_noise.py b/src/ezmsg/simbiophys/line_noise.py index 3b314d5..cf9ecd5 100644 --- a/src/ezmsg/simbiophys/line_noise.py +++ b/src/ezmsg/simbiophys/line_noise.py @@ -89,7 +89,7 @@ class LineNoiseTransformer(BaseStatefulTransformer[LineNoiseSettings, AxisArray, def _chunk_axis(self, message: AxisArray) -> AxisBase | None: """The axis the stream grows along, however the producer named it.""" - dim = message.chunk_dim or next((d for d in self.STREAMING_DIMS if d in message.dims), None) + dim = message.stream_dim or next((d for d in self.STREAMING_DIMS if d in message.dims), None) return message.axes.get(dim) def _hash_message(self, message: AxisArray) -> int: diff --git a/src/ezmsg/simbiophys/noise.py b/src/ezmsg/simbiophys/noise.py index f8c9e9f..d00fd2e 100644 --- a/src/ezmsg/simbiophys/noise.py +++ b/src/ezmsg/simbiophys/noise.py @@ -61,7 +61,7 @@ def _reset_state(self, time_axis: LinearAxis) -> None: dims=["time", "ch"], axes={"time": time_axis, "ch": ch_axis}, # Messages append along `time`; `ch` describes the stream. - chunk_dim="time", + stream_dim="time", ) def _produce(self, n_samples: int, time_axis: LinearAxis) -> AxisArray: @@ -125,7 +125,7 @@ def _initialize_processors(settings: PinkNoiseSettings) -> dict[str, BaseProcess scale=settings.scale, ) ), - # No `axis`: WhiteNoiseProducer declares chunk_dim="time" on what it + # No `axis`: WhiteNoiseProducer declares stream_dim="time" on what it # emits, so the filter resolves to the same dimension it was pinned # to -- without asking ezmsg-sigproc to warn about a setting no user # of PinkNoise ever chose. diff --git a/src/ezmsg/simbiophys/oscillator.py b/src/ezmsg/simbiophys/oscillator.py index 01028bd..d46722f 100644 --- a/src/ezmsg/simbiophys/oscillator.py +++ b/src/ezmsg/simbiophys/oscillator.py @@ -132,7 +132,7 @@ def _reset_state(self, time_axis: LinearAxis) -> None: data=np.zeros((0, 2)), dims=["time", "ch"], axes={"time": time_axis, "ch": ch_axis}, - chunk_dim="time", + stream_dim="time", ) def _produce(self, n_samples: int, time_axis: LinearAxis) -> AxisArray: @@ -242,7 +242,7 @@ def _reset_state(self, time_axis: LinearAxis) -> None: data=np.zeros((0, n_ch)), dims=["time", "ch"], axes={"time": time_axis, "ch": ch_axis}, - chunk_dim="time", + stream_dim="time", ) # Convert settings to arrays and validate diff --git a/tests/unit/test_state_reset_semantics.py b/tests/unit/test_state_reset_semantics.py index 98b856a..5894393 100644 --- a/tests/unit/test_state_reset_semantics.py +++ b/tests/unit/test_state_reset_semantics.py @@ -3,7 +3,7 @@ Two facts these simulators now rely on, neither of which anything else notices if it stops being true: -* Every producer declares :attr:`AxisArray.chunk_dim`, so a downstream processor +* Every producer declares :attr:`AxisArray.stream_dim`, so a downstream processor knows which dimension grows rather than guessing ``"time"``, and hands its channel axis over with the fingerprint already computed, so nobody downstream pays the checksum on every message. @@ -48,7 +48,7 @@ def signal(labels: list[str], fs: float = 100.0, n_time: int = 16, dim: str = "t "ch": CoordinateAxis(data=np.array(labels), dims=["ch"]), }, key="dev", - chunk_dim=dim, + stream_dim=dim, ) @@ -70,10 +70,10 @@ def _produce(self, cls, settings, fs): producer._reset_state(time_axis) return producer._produce(10, time_axis) - def test_declares_the_chunk_dim(self, name, cls, settings, fs): + def test_declares_the_stream_dim(self, name, cls, settings, fs): """Without this a windowing stage downstream has to guess, and ``"time"`` is present-but-wrong the moment the message becomes ``(win, time, ch)``.""" - assert self._produce(cls, settings, fs).chunk_dim == "time" + assert self._produce(cls, settings, fs).stream_dim == "time" def test_hands_over_a_primed_channel_axis(self, name, cls, settings, fs): """The axis is built once per stream, so one checksum covers every @@ -114,7 +114,7 @@ def test_a_relabel_at_a_fixed_count_is_a_new_stream(self, cls, settings): ids=["baseline_drift", "dynamic_colored_noise"], ) def test_a_longer_chunk_is_the_same_stream(self, cls, settings): - """The chunk dimension is excluded, so ordinary chunk-size jitter must + """The stream dimension is excluded, so ordinary chunk-size jitter must not throw away a warmed-up filter.""" proc = cls(settings) assert proc._hash_message(signal(["a", "b"], n_time=16)) == proc._hash_message(signal(["a", "b"], n_time=64)) @@ -132,7 +132,7 @@ def test_a_sample_rate_change_does(self): proc = LineNoiseTransformer(LineNoiseSettings(freq=60.0)) assert proc._hash_message(signal(["a", "b"], fs=100.0)) != proc._hash_message(signal(["a", "b"], fs=500.0)) - def test_it_reads_the_declared_chunk_dim_not_the_name_time(self): + def test_it_reads_the_declared_stream_dim_not_the_name_time(self): """A stream that grows along ``samp`` has no ``time`` axis at all; the old hard-coded lookup silently fell back to a period of zero.""" proc = LineNoiseTransformer(LineNoiseSettings(freq=60.0)) From acc5517c3b1aa1240fe4f4dc154fb9b6c1dffa3d 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 085a43c..c821429 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