Skip to content

mux

BlockMuxStream

BlockMuxStream(reader, start=None, timeout=0, on_gap=ONGAP_DEFAULT)

Bases: Generic[T]

A time-aware, gap-handling multiplexer for SeriesBlock streams.

Given SeriesBlocks from multiple named streams with monotonically increasing integer timestamps, this data structure can be used to pull out sets of synchronized blocks, all with the same timestamps. If data on the streams is not available before timeouts are reached, gap blocks will be returned.

The oldest items will be held until either all named streams are available or until the timeout has been reached. If a start time has been set, any items with an older timestamp will be rejected.

Parameters:

Name Type Description Default
reader StreamReader

StreamReader object producing multiple stream to multiplex.

required
start int

GPS start time of stream, in nanoseconds.

None
timeout int = 0

Overall timeout for the muxer, in nanoseconds. Overrides individual queue timeouts. If not specified the mux timeout will be the max of the individual queue timeouts.

0
on_gap str

Policy for strides in which no stream has any data. 'fill' (default) emits masked gap blocks so the output timeline is continuous; 'skip' emits nothing for such strides (sparse output, no cost for absence); 'raise' raises GapError. Strides where at least one stream has data are always emitted (with the missing channels masked), whatever the policy.

ONGAP_DEFAULT
Source code in arrakis/mux.py
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
def __init__(
    self,
    reader: HasStreams,
    start: int | None = None,
    timeout: int = 0,
    on_gap: str = ONGAP_DEFAULT,
):
    self.reader = reader
    self.start = start
    self.on_gap = OnGap[on_gap.upper()]
    # extract stream info
    self._queues: dict[str, TimedQueue] = {}
    self._stream_channels: dict[str, list[Channel]] = {}
    self.channels: list[Channel] = []
    has_latency_constraint = True
    for stream_name, channels in self.reader.streams.items():
        # gather stride and latency
        # done this way to satisfy type checking consistency
        strides = []
        timeouts = []
        for channel in channels:
            assert channel.stride
            strides.append(channel.stride)
            if channel.max_latency is not None:
                timeouts.append(channel.max_latency)
        # the stride for a stream is the least common multiple of
        # the stride of the individual channels
        qstride = math.lcm(*strides)
        # timeout is max of all expected latencies, or 0 if none
        # have a latency constraint (e.g. historical data)
        if timeouts:
            # one stride of headroom prevents gap-fill from racing
            # data that is in flight between the gRPC thread and
            # the TimedQueue
            qtimeout = max(timeouts) + qstride
        else:
            qtimeout = 0
            has_latency_constraint = False
        self._queues[stream_name] = TimedQueue(
            stride=qstride,
            # mux timeout overrides individual timeouts
            timeout=timeout or qtimeout,
            start=start,
        )
        self._stream_channels[stream_name] = list(channels)
        self.channels.extend(channels)
    # the overall stride for multiple streams is the least common
    # multiple of the stride of the individual streams
    self.stride = math.lcm(*(q.stride for q in self._queues.values()))
    self.timeout = max(q.timeout for q in self._queues.values())
    # Disabled when streams have no latency constraint (e.g.
    # historical data from frames) since there is no wall-clock
    # deadline to gap-fill against.
    self._use_timeouts = has_latency_constraint

__getitem__

__getitem__(key)

Access an individual queue.

Source code in arrakis/mux.py
452
453
454
def __getitem__(self, key: str) -> TimedQueue:
    """Access an individual queue."""
    return self._queues[key]

pull

pull()

Pull synchronized, concatenated, combined blocks from all streams covering the overall specified stride.

Returns:

Type Description
SeriesBlock or None

The combined block for the next stride, or None if the stride held no data and was skipped (on_gap='skip').

Source code in arrakis/mux.py
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
def pull(self) -> SeriesBlock | None:
    """Pull synchronized, concatenated, combined blocks from all
    streams covering the overall specified stride.

    Returns
    -------
    SeriesBlock or None
        The combined block for the next stride, or None if the
        stride held no data and was skipped (``on_gap='skip'``).

    """
    if self.on_gap is not OnGap.FILL and self._next_stride_is_all_gap():
        if self.on_gap is OnGap.RAISE:
            front = min(q.cursor for q in self._queues.values())
            msg = f"gap in stream at {front:_} ns"
            raise GapError(msg)
        self._skip_gap_span()
        return None
    blocks = []
    for stream_name, queue in self._queues.items():
        q_blocks = []
        for time_ns, block in queue.pull(self.stride, update_timeout=False):
            if block is None:
                block = SeriesBlock.full_gap(
                    time_ns,
                    queue.stride,
                    self._stream_channels[stream_name],
                )
            q_blocks.append(block)
        q_block = concatenate_blocks(*q_blocks)
        blocks.append(q_block)
    return combine_blocks(*blocks)

push

push(stream_name, block)

Push an element for time into a particular queue.

Source code in arrakis/mux.py
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
def push(self, stream_name: str, block: SeriesBlock):
    """Push an element for time into a particular queue."""
    q = self._queues[stream_name]
    if block.duration_ns != q.stride:
        # allow partial-stride edge blocks for non-aligned requests
        if q.start is not None and block.duration_ns < q.stride:
            pass
        else:
            logger.warning(
                "dropping block with mismatched stride "
                "at %d ns for stream %s: got %d ns, "
                "expected %d ns",
                block.time_ns,
                stream_name,
                block.duration_ns,
                q.stride,
            )
            return
    q.push(block.time_ns, block)

ready

ready()

True if all queues have the expected number of elements

covering the overall muxer stride.

Source code in arrakis/mux.py
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
def ready(self) -> bool:
    """True if all queues have the expected number of elements

    covering the overall muxer stride.

    """
    # 1. Update timeouts on ALL queues first so gap-filling is
    #    consistent before any alignment happens.
    #    Skip for single-queue muxes — no synchronization needed,
    #    and timeouts race incoming data causing spurious drops.
    if self._use_timeouts:
        for q in self._queues.values():
            q.update_timeout()
        # 2. Align queues to the same front timestamp.
        self._align_queues()
    # 3. Check readiness WITHOUT re-triggering timeouts — the
    #    alignment established in step 2 must not be disturbed.
    return all(
        q.ready(self.stride, update_timeout=False) for q in self._queues.values()
    )

GapError

Bases: Exception

Raised when a stream encounters a gap and on_gap='raise'.

HasStreams

Bases: Protocol

Minimal protocol for objects that provide stream channel metadata.

TimedQueue

TimedQueue(stride, timeout, start=None)

Bases: Generic[T]

A sequential, time-stamped queue handling gaps and timeouts.

The queue stores only real elements; the spans between them are gaps, synthesized lazily when pulled. Two integer cursors define the queue state:

  • cursor: the next timestamp to be served by :meth:pull.
  • horizon: the end (exclusive) of the known span. Every stride in [cursor, horizon) is either a stored element or a gap. Elements older than the horizon are rejected.

The horizon advances when an element is pushed (per-stream delivery is assumed to be in time order, so an element at time T implies nothing older is still coming) and, for live streams, when the timeout deadline passes (absence past the deadline is declared to be a gap). A discontinuity of any size therefore costs no memory: it is served as synthesized gaps between the two stored elements that surround it.

Parameters:

Name Type Description Default
start int

GPS start time of queue, in nanoseconds.

None
stride int

Time step for elements in the queue, in nanoseconds.

required
timeout int

Timeout for elements in the queue, in nanoseconds.

required
Source code in arrakis/mux.py
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
def __init__(self, stride: int, timeout: int, start: int | None = None):
    self.stride = stride
    self.timeout = timeout
    self.start = start
    # Track if queue has received data yet (only relevant when start is None)
    self._initialized = start is not None
    if start is None:
        start = time_as_ns(gpsnow())
    # serve from the stride boundary at or before start, so that a
    # non-aligned edge element at start falls in the first slot
    self.cursor = (start // self.stride) * self.stride
    self.horizon = self.cursor
    # real elements only, strictly increasing times
    self._queue: deque = deque()
    # Backpressure cap on *buffered real elements* (gap spans cost
    # nothing).  Scaled to 2x the timeout window, with a floor of
    # 1000 for historical/no-timeout streams.
    self._max_queue = max(1000, int(2 * self.timeout / self.stride))
    # True once any slot has been served (or the cursor committed
    # by alignment): from then on the serving position must only
    # move forward, so the first-data realignment is disallowed
    self._committed = False
    self._drop_counts: dict[str, int] = {}
    self._drop_log_times: dict[str, float] = {}
    self._lock = RLock()

__len__

__len__()

Number of buffered real elements (gap spans are implicit).

Source code in arrakis/mux.py
155
156
157
def __len__(self):
    """Number of buffered real elements (gap spans are implicit)."""
    return len(self._queue)

drain_until

drain_until(target_time)

Advance the cursor so serving resumes at or after target_time.

Source code in arrakis/mux.py
350
351
352
353
354
355
356
357
358
359
360
361
def drain_until(self, target_time: int) -> None:
    """Advance the cursor so serving resumes at or after target_time."""
    # Snap to stride boundary so the cursor stays stride-aligned.
    target_time = (target_time // self.stride) * self.stride
    with self._lock:
        if target_time <= self.cursor:
            return
        self.cursor = target_time
        self._committed = True
        while self._queue and self._queue[0][0] < self.cursor:
            self._queue.popleft()
        self.horizon = max(self.horizon, self.cursor)

front_time

front_time()

Return the next timestamp to be served, or None if nothing is known.

Source code in arrakis/mux.py
343
344
345
346
347
348
def front_time(self) -> int | None:
    """Return the next timestamp to be served, or None if nothing is known."""
    with self._lock:
        if self.horizon > self.cursor or self._queue:
            return self.cursor
        return None

next_element_time

next_element_time()

Return the timestamp of the oldest buffered element, or None.

Unlike :meth:front_time this ignores gap spans: it reports where the next real element is, however far ahead.

Source code in arrakis/mux.py
332
333
334
335
336
337
338
339
340
341
def next_element_time(self) -> int | None:
    """Return the timestamp of the oldest buffered element, or None.

    Unlike :meth:`front_time` this ignores gap spans: it reports
    where the next real element is, however far ahead.
    """
    with self._lock:
        if self._queue:
            return self._queue[0][0]
        return None

pull

pull(duration=None, *, update_timeout=True)

Drain the queue.

Gaps are represented by None elements, synthesized on the fly.

Parameters:

Name Type Description Default
duration int

Duration to extract, in nanoseconds. If the specified duration is not available, no elements will be returned. If not specified, the known span will be drained.

None
update_timeout bool

Whether to trigger timeout gap-filling before pulling. Default is True. Set to False when the caller has already updated timeouts (e.g. BlockMuxStream.pull()).

True

Yields:

Type Description
tuple[time, element]

The element from the queue and it's associated timestamp.

Source code in arrakis/mux.py
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
def pull(
    self, duration: int | None = None, *, update_timeout: bool = True
) -> Iterator[tuple[int, Any]]:
    """Drain the queue.

    Gaps are represented by None elements, synthesized on the fly.

    Parameters
    ----------
    duration : int
        Duration to extract, in nanoseconds.  If the specified
        duration is not available, no elements will be returned.
        If not specified, the known span will be drained.
    update_timeout : bool, optional
        Whether to trigger timeout gap-filling before pulling.
        Default is True.  Set to False when the caller has already
        updated timeouts (e.g. BlockMuxStream.pull()).

    Yields
    ------
    tuple[time, element]
        The element from the queue and it's associated timestamp.

    """
    if update_timeout:
        self.update_timeout()

    with self._lock:
        # if duration specified, pull the requested number of elements
        if duration:
            if n_elements := self.ready(duration, update_timeout=False):
                for _ in range(n_elements):
                    yield self._serve()

        # else drain the known span
        else:
            while self.cursor < self.horizon or self._queue:
                yield self._serve()

push

push(time, element, on_drop=ONDROP_DEFAULT)

Push an element into the queue.

The time being pushed into the queue must be a multiple of the time stride specified at initialization of the queue.

If the time associated with the pushed element is older than the queue's horizon (already served or declared gap), the element will be dropped and this operation will be a no-op.

Pushing None declares the span up to time to be a gap without storing anything.

Parameters:

Name Type Description Default
time int

GPS time associated with the element, in nanoseconds

required
element Any

element being pushed into the queue.

required
on_drop str

Per-event behavior when the item is dropped as too old (e.g. it arrived after the timeout declared its span a gap, or it is a duplicate). Options are 'ignore', 'raise', or 'warn'. Default is 'ignore': drops are expected under the latency contract, and are always counted and summarized in the log at most once per DROP_LOG_INTERVAL regardless of this policy.

ONDROP_DEFAULT
Source code in arrakis/mux.py
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
def push(self, time: int, element: Any, on_drop: str = ONDROP_DEFAULT) -> None:
    """Push an element into the queue.

    The time being pushed into the queue must be a multiple of the
    time stride specified at initialization of the queue.

    If the time associated with the pushed element is older than
    the queue's horizon (already served or declared gap), the
    element will be dropped and this operation will be a no-op.

    Pushing ``None`` declares the span up to ``time`` to be a gap
    without storing anything.

    Parameters
    ----------
    time : int
        GPS time associated with the element, in nanoseconds
    element : Any
        element being pushed into the queue.
    on_drop : str, optional
        Per-event behavior when the item is dropped as too old
        (e.g. it arrived after the timeout declared its span a
        gap, or it is a duplicate).  Options are 'ignore', 'raise',
        or 'warn'.  Default is 'ignore': drops are expected under
        the latency contract, and are always counted and
        summarized in the log at most once per DROP_LOG_INTERVAL
        regardless of this policy.

    """
    assert time % self.stride == 0 or time == self.start, (
        f"time {time} is not a multiple of queue stride {self.stride:_}"
        f" (and does not match start {self.start})"
    )
    # If this is the first real element on a live queue, realign to
    # it: serving starts at its stride boundary, and any span the
    # timeout may have declared gap in the meantime is discarded.
    # Only allowed while nothing has been served yet — once a slot
    # has been yielded (e.g. a timeout gap emitted just before the
    # first element landed), rewinding would duplicate it, so the
    # element falls through to the normal too-old handling instead.
    if not self._initialized and element is not None:
        if not self._committed:
            with self._lock:
                self.cursor = (time // self.stride) * self.stride
                self.horizon = self.cursor
        self._initialized = True

    # if time is older than the horizon, drop it
    if time < self.horizon:
        msg = f"item's timestamp is too old: ({time:_} < {self.horizon:_})"
        match OnDrop[on_drop.upper()]:
            case OnDrop.IGNORE:
                if element is not None:
                    self._count_drop("late", msg)
                return
            case OnDrop.RAISE:
                raise ValueError(msg)
            case OnDrop.WARN:
                if element is not None:
                    self._count_drop("late", msg)
                logger.warning(msg)
                warnings.warn(msg, stacklevel=2)
                return
    with self._lock:
        if element is not None:
            # Never raise on overflow: this runs inside long-lived
            # server poll threads, where an exception silently
            # kills the stream for every subscriber.  Bound memory
            # by dropping the oldest element instead; the dropped
            # span is served as a gap.
            if len(self._queue) >= self._max_queue:
                self._queue.popleft()
                self._count_drop(
                    "overflow",
                    f"queue at capacity ({self._max_queue} elements)",
                )
            self._queue.append((time, element))
        # the element's slot is now known; for a non-aligned edge
        # element this is the enclosing stride boundary
        self.horizon = (time // self.stride) * self.stride + self.stride

ready

ready(duration, *, update_timeout=True)

Check if queue holds duration worth of elements

Parameters:

Name Type Description Default
duration int

Duration to check for, in nanoseconds.

required
update_timeout bool

Whether to trigger timeout gap-filling before checking. Default is True. Set to False when the caller has already updated timeouts (e.g. BlockMuxStream.ready()).

True

Returns:

Type Description
int or None

Returns either the number of elements that span duration, or None if the known span does not cover duration.

Source code in arrakis/mux.py
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
def ready(self, duration: int, *, update_timeout: bool = True) -> int | None:
    """Check if queue holds duration worth of elements

    Parameters
    ----------
    duration : int
        Duration to check for, in nanoseconds.
    update_timeout : bool, optional
        Whether to trigger timeout gap-filling before checking.
        Default is True.  Set to False when the caller has already
        updated timeouts (e.g. BlockMuxStream.ready()).

    Returns
    -------
    int or None
        Returns either the number of elements that span duration,
        or None if the known span does not cover duration.

    """
    if update_timeout:
        self.update_timeout()
    assert duration % self.stride == 0, (
        f"duration {duration:_} is not a multiple of queue stride {self.stride:_}"
    )
    if self.horizon - self.cursor >= duration:
        return int(duration / self.stride)
    return None