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
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

    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

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

    @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"

        if all(frame.is_gap for frame in frames):
            # Return a gap buffer if all frames are gaps
            out = None
            shape = frames[0].shape
        else:
            # A reference array from any non-gap buffer fixes the backend, dtype,
            # and device for any zeros we fill gaps with (the array is the
            # backend). Guaranteed present here since not all frames are gaps.
            ref = next(buf.data for f in frames for buf in f if not buf.is_gap)
            xp = array_namespace(ref)
            assert xp is not None

            # use the first frame as basis
            if len(frames[0]) == 1:
                out = frames[0][0].filleddata(ref)
            else:
                out = xp.concat([buf.filleddata(ref) for buf in frames[0]], axis=-1)
            shape = out.shape
            # add to the first frame
            for i, f in enumerate(frames[1:]):
                i0 = 0
                for buf in f:
                    if not buf.is_gap:
                        if self.addslices_map is None:
                            out[..., i0 : i0 + buf.samples] += buf.data
                        else:
                            slices = self.addslices_map[self.sink_pad_names[i + 1]] + (
                                slice(i0, i0 + buf.samples),
                            )
                            out[slices] += buf.data

                    i0 += buf.samples

        output_frame.append(
            SeriesBuffer(
                offset=frames[0].offset,
                sample_rate=frames[0].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"

    if all(frame.is_gap for frame in frames):
        # Return a gap buffer if all frames are gaps
        out = None
        shape = frames[0].shape
    else:
        # A reference array from any non-gap buffer fixes the backend, dtype,
        # and device for any zeros we fill gaps with (the array is the
        # backend). Guaranteed present here since not all frames are gaps.
        ref = next(buf.data for f in frames for buf in f if not buf.is_gap)
        xp = array_namespace(ref)
        assert xp is not None

        # use the first frame as basis
        if len(frames[0]) == 1:
            out = frames[0][0].filleddata(ref)
        else:
            out = xp.concat([buf.filleddata(ref) for buf in frames[0]], axis=-1)
        shape = out.shape
        # add to the first frame
        for i, f in enumerate(frames[1:]):
            i0 = 0
            for buf in f:
                if not buf.is_gap:
                    if self.addslices_map is None:
                        out[..., i0 : i0 + buf.samples] += buf.data
                    else:
                        slices = self.addslices_map[self.sink_pad_names[i + 1]] + (
                            slice(i0, i0 + buf.samples),
                        )
                        out[slices] += buf.data

                i0 += buf.samples

    output_frame.append(
        SeriesBuffer(
            offset=frames[0].offset,
            sample_rate=frames[0].sample_rate,
            data=out,
            shape=shape,
        )
    )