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
41 changes: 31 additions & 10 deletions src/ezmsg/sigproc/resample.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,14 +22,30 @@
from .util.message import has_samples_along


def _as_limit(value: float | None) -> float | None:
"""Normalize an optional "no limit" setting to ``None``.

``None`` is the canonical way to say "no limit" because it survives JSON
serialization, which ``inf`` does not. Non-finite values are accepted and
folded onto ``None`` so pipelines that predate this convention keep working.
"""
if value is None or not math.isfinite(value):
return None
return value


class ResampleSettings(ez.Settings):
axis: str = "time"

resample_rate: float | None = None
"""target resample rate in Hz. If None, the resample rate will be determined by the reference signal."""

max_chunk_delay: float = np.inf
"""Maximum delay between outputs in seconds. If the delay exceeds this value, the transformer will extrapolate."""
max_chunk_delay: float | None = None
"""
Maximum delay between outputs in seconds. If the delay exceeds this value, the
transformer will extrapolate. ``None`` (the default) means no limit: the transformer
never extrapolates on wall-clock alone. A non-finite value is treated as ``None``.
"""

fill_value: str = "extrapolate"
"""
Expand Down Expand Up @@ -62,7 +78,7 @@ class ResampleSettings(ez.Settings):
rate mode the reference grid is synthetic and carries no data.
"""

reference_reset_after_chunks: float = 3
reference_reset_after_chunks: int | None = 3
"""
Robustness against a non-monotonic reference clock.

Expand All @@ -78,9 +94,10 @@ class ResampleSettings(ez.Settings):
resumes on the new clock. Small, self-correcting jitter (a few out-of-order samples)
does not trigger this and is simply skipped, keeping the output monotonic.

Set to ``float("inf")`` to disable reset recovery (the transformer may then stall
Set to ``None`` to disable reset recovery (the transformer may then stall
indefinitely on a backward clock jump). Output across a genuine reset is necessarily
discontinuous; sanitising the reference timestamps upstream remains the robust fix.
A non-finite value is treated as ``None``.
"""


Expand Down Expand Up @@ -225,7 +242,8 @@ def push_reference(self, message: AxisArray) -> None:
self.state.stale_ref_pushes += 1
else:
self.state.stale_ref_pushes = 0
if self.state.stale_ref_pushes >= self.settings.reference_reset_after_chunks:
reset_after = _as_limit(self.settings.reference_reset_after_chunks)
if reset_after is not None and self.state.stale_ref_pushes >= reset_after:
first, step = self._axis_first_step(ax)
warnings.warn(
"ResampleProcessor: reference clock jumped backward and stayed behind "
Expand Down Expand Up @@ -297,11 +315,14 @@ def __next__(self) -> AxisArray:

# If we do not rely on an external reference, and we have not received new data in a while,
# then extrapolate our reference vector out beyond the delay limit.
b_project = self.settings.resample_rate is not None and time.monotonic() > (
self.state.last_write_time + self.settings.max_chunk_delay
max_chunk_delay = _as_limit(self.settings.max_chunk_delay)
b_project = (
self.settings.resample_rate is not None
and max_chunk_delay is not None
and time.monotonic() > (self.state.last_write_time + max_chunk_delay)
)
if b_project:
n_append = math.ceil(self.settings.max_chunk_delay / ref_ax.gain)
n_append = math.ceil(max_chunk_delay / ref_ax.gain)
xvec_append = ref_xvec[-1] + np.arange(1, n_append + 1) * ref_ax.gain
ref_xvec = np.hstack((ref_xvec, xvec_append))

Expand Down Expand Up @@ -449,8 +470,8 @@ async def gen_resampled(self):
# the wall-clock extrapolation in `__next__` can fire without input.
# Reference-driven mode never becomes ready by wall-clock, so there
# a pure event wait suffices.
timeout = self.SETTINGS.max_chunk_delay if self.SETTINGS.resample_rate is not None else np.inf
if np.isfinite(timeout):
timeout = _as_limit(self.SETTINGS.max_chunk_delay) if self.SETTINGS.resample_rate is not None else None
if timeout is not None:
try:
await asyncio.wait_for(self._wake.wait(), timeout=timeout)
except asyncio.TimeoutError:
Expand Down
6 changes: 3 additions & 3 deletions src/ezmsg/sigproc/resampleconcat.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,6 @@
import typing

import ezmsg.core as ez
import numpy as np
from ezmsg.util.messages.axisarray import AxisArray

from .concat import ConcatProcessor, ConcatSettings
Expand All @@ -48,9 +47,10 @@ class ResampleConcatSettings(ez.Settings):
# --- Resample passthrough ---
buffer_duration: float = 2.0
fill_value: str = "extrapolate"
max_chunk_delay: float = np.inf
max_chunk_delay: float | None = None
"""See :attr:`ezmsg.sigproc.resample.ResampleSettings.max_chunk_delay`."""
buffer_update_strategy: UpdateStrategy = "immediate"
reference_reset_after_chunks: float = 3
reference_reset_after_chunks: int | None = 3
"""See :attr:`ezmsg.sigproc.resample.ResampleSettings.reference_reset_after_chunks`."""

# --- Concat passthrough ---
Expand Down
4 changes: 2 additions & 2 deletions tests/integration/ezmsg/test_resample_merge_system.py
Original file line number Diff line number Diff line change
Expand Up @@ -214,14 +214,14 @@ def test_resample_merge_healthy_system(test_name: str | None = None):
def test_resample_merge_seizes_when_recovery_disabled_system(test_name: str | None = None):
"""Documents the original bug: with recovery disabled, a backward reference jump stops output.

Setting ``reference_reset_after_chunks=inf`` restores the pre-hardening behaviour, so a
Setting ``reference_reset_after_chunks=None`` restores the pre-hardening behaviour, so a
large sustained backward offset jump pushes the reference permanently below the
resampler's high-water mark and the merged output ceases well before input is exhausted.
"""
_, _, _, last_t = _run_resample_merge(
glitch_at=15,
glitch_back=100.0,
reset_after=float("inf"),
reset_after=None,
use_output_reference=False,
fn=get_test_fn(test_name),
)
Expand Down
11 changes: 8 additions & 3 deletions tests/unit/test_resample.py
Original file line number Diff line number Diff line change
Expand Up @@ -250,10 +250,15 @@ def test_resample_recovers_from_reference_reset():
assert sum(counts[14:]) > 0, "Resampler never recovered after the reference reset."


def test_resample_reset_disabled_can_stall():
"""With recovery disabled (inf threshold), a large backward jump stops output."""
@pytest.mark.parametrize("reset_after", [None, float("inf")], ids=["none", "legacy_inf"])
def test_resample_reset_disabled_can_stall(reset_after):
"""With recovery disabled, a large backward jump stops output.

``None`` is the canonical "disabled"; the legacy non-finite spelling is folded
onto it so pre-existing pipelines keep working (``inf`` does not survive JSON).
"""
fs, chunk = 100.0, 30
resample = ResampleProcessor(resample_rate=None, buffer_duration=4.0, reference_reset_after_chunks=float("inf"))
resample = ResampleProcessor(resample_rate=None, buffer_duration=4.0, reference_reset_after_chunks=reset_after)
counts = []
for i in range(25):
off = i * chunk / fs - (100.0 if i >= 10 else 0.0)
Expand Down
Loading