ezmsg.baseproc.util.streamdim#

Which dimension a processor should operate on.

An AxisArray names its dimensions but nothing about a position says what a dimension means. dims[0] is not “the streaming axis” and dims[-1] is not “the channel axis”; both guesses break under transpose() and downstream of any windowing stage, where a (time, ch) stream becomes (win, time, ch).

stream_dim is the producer’s declaration of which dimension messages accumulate along – the one party that reliably knows. These helpers turn that declaration into the axis a given kind of processor should use, and they live here because _message_hash() already resolves the same thing for its own purposes: a processor whose arithmetic disagreed with its state-reset logic about which dimension is which would reset on the wrong changes and cache state along the wrong axis.

Three rules, because one does not fit every case:

Module Attributes

STREAMING_DIMS

Default fallback stream dimension, matching BaseStatefulTransformer.

Functions

has_samples_along(message, dim)[source]#

True iff dim is present in message with nonzero length.

Stricter than not is_empty_along(...): the dim must exist. Drain loops use this to decide whether a chunk is real output, so that a placeholder lacking the axis entirely (e.g. ResampleProcessor’s pre-init null template, dims=[""]) counts as “nothing ready” rather than a publishable chunk.

Return type:

bool

Parameters:
is_empty_along(message, dims)[source]#

True iff any of the named dims is present in message with zero length.

Publish gates use this instead of data.size == 0 so a message that is empty only along other axes — e.g. an upstream selection removed every channel while time samples remain — still flows downstream, preserving the stream’s cadence for consumers that align or merge multiple sources. Dims not present in the message are ignored.

Return type:

bool

Parameters:
resolve_configured_stream_dim(processor, message, configured, legacy_default=None)[source]#

Resolve a state-carrying processor’s axis, honouring an explicit setting.

configured wins when set – an explicit axis is an instruction, and removing that escape hatch would break every pipeline that passes the common axis="time". But when the producer declared a different stream dimension, that disagreement is worth surfacing exactly once: the processor’s cross-message state is about to be carried along an axis whose length is fixed, which is a different operation from the one the caller almost certainly meant.

The warning fires only against a declared stream_dim, never against the STREAMING_DIMS guess – warning on a guess would fire on every correctly-configured windowed pipeline whose producer is merely silent.

Parameters:
  • legacy_default (str | None) – The dimension this processor’s axis setting used to default to, for the stages whose default was a hardcoded "time" rather than a positional guess. Flipping those to follow stream_dim changes results wherever the stream dimension is not "time" – most obviously downstream of a windowing stage, where it is "win" – and unlike an explicitly configured axis there is nothing in the settings to warn about. Passing the old default here surfaces exactly that population, once, and is dropped when the setting is removed.

  • processor (Any)

  • message (AxisArray)

  • configured (str | None)

Return type:

str

resolve_feature_dim(message, position=-1)[source]#

The dimension at position, skipping the stream dimension.

For processors whose axis is a static one – channels, coordinate components, feature labels. stream_dim is emphatically not the answer here, but the naive dims[position] can silently be the stream dimension: a (ch, time) stream makes dims[-1] the accumulating axis, and an affine transform would then matmul across time while a slicer would discard samples.

Falls back to dims[position] when the stream dimension is all there is, which keeps 1-D messages working rather than raising on them.

Return type:

str

Parameters:
resolve_stream_dim(message, streaming_dims=('time',))[source]#

The dimension successive messages accumulate along.

This is the axis a processor that carries state between messages must operate on – filter initial conditions, a running mean, a sample buffer, a previous-sample cache. Carrying such state along any other dimension is not a smaller error but a different operation: a static axis has the same length every message, so state carried across it applies message N’s tail to message N+1’s head at the same coordinate, forever.

The producer renamed the dims and so is the only party that reliably knows which one grows; message.stream_dim is that declaration. When a producer is silent, streaming_dims supplies the guess – ("time",) is right for a raw signal and wrong downstream of a windowing stage, where the message is (win, time, ch) and win is what grows.

dims[0] is the last resort only. It is a position, not a meaning, and it breaks under transpose().

Return type:

str

Parameters:
resolve_transform_dim(message, streaming_dims=('time',))[source]#

The regularly-sampled dimension a transform consumes.

Neither resolve_stream_dim() nor resolve_feature_dim() fits a stage like Spectrum, which needs the axis whose gain is a sample period and whose extent is the transform length:

  • On a raw (time, ch) stream that is the stream dimension.

  • On windowed (win, time, ch) it is timewin is what accumulates, but each window’s spectrum is taken over time.

So: prefer the innermost non-stream dimension carrying a LinearAxis, and fall back to the stream dimension when there is none. ch carries a CoordinateAxis (or no axis at all), so the raw case falls through correctly rather than transforming across channels.

Return type:

str

Parameters:
with_fingerprint(axis)[source]#

Compute axis’s fingerprint now, and return the axis.

Every stateful consumer reads the fingerprint of the coordinate axes that describe a stream’s configuration, and the value is cached on the instance and pickled with it. Computing it at the point of construction therefore pays the checksum once, for everybody:

  • In this process, the axis object is reused for the life of the stream, so one call covers every message and every consumer downstream of it.

  • Across a process boundary it is better than that. Unpickling hands out a new axis object per message, so a cold axis is re-checksummed by the first consumer in every receiving process, on every message, forever. A primed one arrives with the answer already attached.

Apply it to axes that describe the stream – channel labels, frequency labels, feature labels – not to per-message coordinates along the stream dimension, whose fingerprint no consumer reads and whose data is new every message anyway.

Return type:

CoordinateAxis

Parameters:

axis (CoordinateAxis)