Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions .github/workflows/python-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -19,6 +25,7 @@ jobs:
- "windows-latest"
- "macos-latest"
runs-on: ${{matrix.os}}
timeout-minutes: 30

steps:
- uses: actions/checkout@v4
Expand Down
4 changes: 2 additions & 2 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
2 changes: 1 addition & 1 deletion src/ezmsg/xdf/iter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
)


Expand Down
18 changes: 9 additions & 9 deletions tests/test_iter.py
Original file line number Diff line number Diff line change
Expand Up @@ -109,21 +109,21 @@ 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
on the axis and pickled with it; priming it at construction spares the first
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]
Expand All @@ -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"
Expand All @@ -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

Expand Down Expand Up @@ -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):
Expand All @@ -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)
Loading