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:
resolve_stream_dim()– for state carried between messages.resolve_feature_dim()– for a static axis (channels, components).resolve_transform_dim()– for a transform that consumes a regularly sampled axis, which downstream of a windowing stage is not the stream one.
Module Attributes
Default fallback stream dimension, matching |
Functions
- has_samples_along(message, dim)[source]#
True iff
dimis present inmessagewith 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.
- is_empty_along(message, dims)[source]#
True iff any of the named dims is present in
messagewith zero length.Publish gates use this instead of
data.size == 0so 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.
- 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 theSTREAMING_DIMSguess – 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’saxissetting used to default to, for the stages whose default was a hardcoded"time"rather than a positional guess. Flipping those to followstream_dimchanges 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:
- 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_dimis emphatically not the answer here, but the naivedims[position]can silently be the stream dimension: a(ch, time)stream makesdims[-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.
- 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_dimis 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)andwinis what grows.dims[0]is the last resort only. It is a position, not a meaning, and it breaks undertranspose().
- resolve_transform_dim(message, streaming_dims=('time',))[source]#
The regularly-sampled dimension a transform consumes.
Neither
resolve_stream_dim()norresolve_feature_dim()fits a stage likeSpectrum, which needs the axis whosegainis 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 istime–winis what accumulates, but each window’s spectrum is taken overtime.
So: prefer the innermost non-stream dimension carrying a
LinearAxis, and fall back to the stream dimension when there is none.chcarries aCoordinateAxis(or no axis at all), so the raw case falls through correctly rather than transforming across channels.
- 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:
- Parameters:
axis (CoordinateAxis)