Skip to content

sgnts.transforms.adder

Adder dataclass

Bases: TSTransform


              flowchart TD
              sgnts.transforms.adder.Adder[Adder]
              sgnts.base.base.TSTransform[TSTransform]
              sgnts.base.base.TimeSeriesMixin[TimeSeriesMixin]

                              sgnts.base.base.TSTransform --> sgnts.transforms.adder.Adder
                                sgnts.base.base.TimeSeriesMixin --> sgnts.base.base.TSTransform
                



              click sgnts.transforms.adder.Adder href "" "sgnts.transforms.adder.Adder"
              click sgnts.base.base.TSTransform href "" "sgnts.base.base.TSTransform"
              click sgnts.base.base.TimeSeriesMixin href "" "sgnts.base.base.TimeSeriesMixin"
            

Add up all the frames from all the sink pads.

Parameters:

Name Type Description Default
addslices_map dict[str, tuple[slice, ...]] | None

Optional[dict[str, tuple[slice, ...]], a mapping of sink_pad_names to a tuple of slice objects, representing array index slices in each dimension except the last. Suppose there are two sink pads "sink_pad_name1" and "sink_pad_name2", and data1 is the data from sink_pad_name1, and data2 is the data from sink_pad_name2, and addslices_map = {"sink_pad_name2": (slice(2, 6), slice(0, 8))}, then this element will perform the following operation:

out = data1[slice(2, 6), slice(0, 8), :] + data2
None
fill_gaps bool

bool, if True (the default), treat input gap buffers as zeros when summing, only outputting a gap where all inputs have gaps. If False, output a gap wherever any input has a gap.

True
Notes

Thread safety: Marked thread_safe = True. With Pipeline.run(threaded=N) the pad callbacks for this element are dispatched onto worker threads.

Pad layout: N sink pads + 1 source pad
(``@validator.many_to_one``). The N sink pads' ``pull``
callbacks CAN run concurrently in the same wave — that is
the per-pad concurrency to keep safe. ``internal`` runs
alone.

Where the GIL-releasing work lives: ``internal()`` →
``process()`` does NumPy/Torch element-wise addition,
which releases the GIL for large arrays.

Per-pad concurrency analysis:

- ``pull`` (inherited ``TimeSeriesMixin.pull``): writes
  per-pad-keyed dicts (``inbufs[pad]``, ``metadata[pad]``).
  Distinct keys per pad → safe under concurrent calls.
  Also OR's ``self.at_EOS``, which is idempotent for
  booleans (any True wins).
- ``new`` (inherited): read-only lookup in
  ``self.outframes``.
- ``process``: reads each input frame and accumulates into
  a local ``out`` array. No element-level mutation.

**Future editors MUST preserve thread safety**: keep
``process`` purely functional on its inputs and a local
output. Do NOT add element-level state mutated from
``pull`` outside of per-pad-keyed containers — multiple
sink pads will race on it under threading.
Source code in src/sgnts/transforms/adder.py
@dataclass
class Adder(TSTransform):
    """Add up all the frames from all the sink pads.

    Args:
        addslices_map:
            Optional[dict[str, tuple[slice, ...]], a mapping of sink_pad_names to a
            tuple of slice objects, representing array index slices in each dimension
            except the last. Suppose there are two sink pads "sink_pad_name1" and
            "sink_pad_name2", and data1 is the data from sink_pad_name1, and data2 is
            the data from sink_pad_name2, and addslices_map = {"sink_pad_name2":
            (slice(2, 6), slice(0, 8))}, then this element will perform the following
            operation:

                out = data1[slice(2, 6), slice(0, 8), :] + data2
        fill_gaps:
            bool, if True (the default), treat input gap buffers as zeros when
            summing, only outputting a gap where all inputs have gaps. If
            False, output a gap wherever any input has a gap.

    Notes:
        Thread safety:
            Marked ``thread_safe = True``. With
            ``Pipeline.run(threaded=N)`` the pad callbacks for this
            element are dispatched onto worker threads.

            Pad layout: N sink pads + 1 source pad
            (``@validator.many_to_one``). The N sink pads' ``pull``
            callbacks CAN run concurrently in the same wave — that is
            the per-pad concurrency to keep safe. ``internal`` runs
            alone.

            Where the GIL-releasing work lives: ``internal()`` →
            ``process()`` does NumPy/Torch element-wise addition,
            which releases the GIL for large arrays.

            Per-pad concurrency analysis:

            - ``pull`` (inherited ``TimeSeriesMixin.pull``): writes
              per-pad-keyed dicts (``inbufs[pad]``, ``metadata[pad]``).
              Distinct keys per pad → safe under concurrent calls.
              Also OR's ``self.at_EOS``, which is idempotent for
              booleans (any True wins).
            - ``new`` (inherited): read-only lookup in
              ``self.outframes``.
            - ``process``: reads each input frame and accumulates into
              a local ``out`` array. No element-level mutation.

            **Future editors MUST preserve thread safety**: keep
            ``process`` purely functional on its inputs and a local
            output. Do NOT add element-level state mutated from
            ``pull`` outside of per-pad-keyed containers — multiple
            sink pads will race on it under threading.
    """

    thread_safe = True

    # concat / zeros are standard-xp ops; works in any namespace.
    backends = ANY_BACKEND

    addslices_map: dict[str, tuple[slice, ...]] | None = None
    fill_gaps: bool = True

    def configure(self) -> None:
        self.adapter_config.alignment(
            align_buffers=True,
        )

    @validator.many_to_one
    def validate(self) -> None:
        pass

    def output_prototype(self, pad):
        return sum(self.input_prototype(spad.pad_name) for spad in self.sink_pads)

    @transform.many_to_one
    def process(
        self, input_frames: dict[SinkPad, TSFrame], output_frame: TSCollectFrame
    ) -> None:
        """Add up all the frames from all the sink pads."""
        frames = list(input_frames.values())

        # Sanity check frames
        assert (
            len({f.sample_rate for f in frames}) == 1
        ), "Sample rate of frames must be the same"
        assert len({f.offset for f in frames}) == 1, "Frames must be aligned"
        assert len({f.end_offset for f in frames}) == 1, "Frames must be aligned"

        if self.addslices_map is None:
            assert (
                len({f.shape for f in frames}) == 1
            ), "Shape of frames must be the same"
        else:
            assert (
                len({f.shape[-1] for f in frames}) == 1
            ), "Size of last dimension must be the same"

        sample_rate = frames[0].sample_rate

        if all(frame.is_gap for frame in frames):
            # Return a single gap buffer if all frames are gaps
            output_frame.append(
                SeriesBuffer(
                    offset=frames[0].offset,
                    sample_rate=sample_rate,
                    data=None,
                    shape=frames[0].shape,
                )
            )
            return

        proto = self.output_prototype(self.source_pads[0])
        assert proto is not None

        # iterate over aligned buffers from all input frames
        for bufs in zip(*frames):
            offset = bufs[0].offset
            shape = bufs[0].shape

            # output a gap if all the buffers are gaps, or, when not
            # filling gaps, if any of the buffers is a gap
            if (
                all(buf.is_gap for buf in bufs)
                if self.fill_gaps
                else any(buf.is_gap for buf in bufs)
            ):
                out = None

            # otherwise sum the data from all the buffers, assuming
            # gaps are zero
            else:
                # follow the prototype's namespace, dtype *and* device so the
                # accumulator lands on the same device as the incoming data
                # (plain xp.zeros defaults torch tensors to CPU)
                out = new_zeros(proto, shape)

                for pad_name, buf in zip(self.sink_pad_names, bufs):
                    if buf.is_gap:
                        continue
                    if self.addslices_map is None or pad_name not in self.addslices_map:
                        out += buf.data
                    else:
                        slices = self.addslices_map[pad_name] + (slice(0, buf.samples),)
                        out[slices] += buf.data

            output_frame.append(
                SeriesBuffer(
                    offset=offset,
                    sample_rate=sample_rate,
                    data=out,
                    shape=shape,
                )
            )

process(input_frames, output_frame)

Add up all the frames from all the sink pads.

Source code in src/sgnts/transforms/adder.py
@transform.many_to_one
def process(
    self, input_frames: dict[SinkPad, TSFrame], output_frame: TSCollectFrame
) -> None:
    """Add up all the frames from all the sink pads."""
    frames = list(input_frames.values())

    # Sanity check frames
    assert (
        len({f.sample_rate for f in frames}) == 1
    ), "Sample rate of frames must be the same"
    assert len({f.offset for f in frames}) == 1, "Frames must be aligned"
    assert len({f.end_offset for f in frames}) == 1, "Frames must be aligned"

    if self.addslices_map is None:
        assert (
            len({f.shape for f in frames}) == 1
        ), "Shape of frames must be the same"
    else:
        assert (
            len({f.shape[-1] for f in frames}) == 1
        ), "Size of last dimension must be the same"

    sample_rate = frames[0].sample_rate

    if all(frame.is_gap for frame in frames):
        # Return a single gap buffer if all frames are gaps
        output_frame.append(
            SeriesBuffer(
                offset=frames[0].offset,
                sample_rate=sample_rate,
                data=None,
                shape=frames[0].shape,
            )
        )
        return

    proto = self.output_prototype(self.source_pads[0])
    assert proto is not None

    # iterate over aligned buffers from all input frames
    for bufs in zip(*frames):
        offset = bufs[0].offset
        shape = bufs[0].shape

        # output a gap if all the buffers are gaps, or, when not
        # filling gaps, if any of the buffers is a gap
        if (
            all(buf.is_gap for buf in bufs)
            if self.fill_gaps
            else any(buf.is_gap for buf in bufs)
        ):
            out = None

        # otherwise sum the data from all the buffers, assuming
        # gaps are zero
        else:
            # follow the prototype's namespace, dtype *and* device so the
            # accumulator lands on the same device as the incoming data
            # (plain xp.zeros defaults torch tensors to CPU)
            out = new_zeros(proto, shape)

            for pad_name, buf in zip(self.sink_pad_names, bufs):
                if buf.is_gap:
                    continue
                if self.addslices_map is None or pad_name not in self.addslices_map:
                    out += buf.data
                else:
                    slices = self.addslices_map[pad_name] + (slice(0, buf.samples),)
                    out[slices] += buf.data

        output_frame.append(
            SeriesBuffer(
                offset=offset,
                sample_rate=sample_rate,
                data=out,
                shape=shape,
            )
        )