Skip to content

client

Client-based API access.

Client

Client(url=None, *, token=None)

Retrieve channel information and timeseries data from Arrakis.

Parameters:

Name Type Description Default
url str

The URL to connect to. If the URL is not set, connect to a default server or one set by ARRAKIS_SERVER.

None
token str, bool, or None

Controls authentication. None (default) auto-discovers a token via igwn-auth-utils, falling back to unauthenticated. True auto-discovers but raises if no token is found. False disables authentication. A string is used as the raw JWT directly.

None
Source code in arrakis/client.py
148
149
150
151
152
153
154
155
156
157
158
def __init__(
    self,
    url: str | None = None,
    *,
    token: str | bool | None = None,
):
    self.initial_url = parse_arrakis_url(url).geturl()
    logger.debug("initial url: %s", self.initial_url)
    self._token = resolve_token(token, self.initial_url)
    self._middleware = build_auth_middleware(self._token)
    self._validator = RequestValidator()

count

count(pattern=constants.DEFAULT_MATCH, data_type=None, min_rate=constants.MIN_SAMPLE_RATE, max_rate=constants.MAX_SAMPLE_RATE, publisher=None, time=None, replay_id=None)

Count channels matching a set of conditions

Parameters:

Name Type Description Default
pattern str

Channel pattern to match channels with, using regular expressions.

DEFAULT_MATCH
data_type dtype - like | list[dtype - like]

If set, find all channels with these data types.

None
min_rate int

The minimum sampling rate for channels.

MIN_SAMPLE_RATE
max_rate int

The maximum sampling rate for channels.

MAX_SAMPLE_RATE
publisher str | list[str]

If set, find all channels associated with these publishers.

None
time float

GPS time in seconds indicating when the metadata query is valid. If None, routes to the live backend (current state).

None
replay_id str

Server-registered replay identifier.

None

Returns:

Type Description
int

The number of channels matching query.

Source code in arrakis/client.py
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
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
279
280
281
282
283
284
285
286
287
def count(
    self,
    pattern: str = constants.DEFAULT_MATCH,
    data_type: DataTypeLike | None = None,
    min_rate: int | None = constants.MIN_SAMPLE_RATE,
    max_rate: int | None = constants.MAX_SAMPLE_RATE,
    publisher: str | list[str] | None = None,
    time: float | None = None,
    replay_id: str | None = None,
) -> int:
    """Count channels matching a set of conditions

    Parameters
    ----------
    pattern : str, optional
        Channel pattern to match channels with, using regular expressions.
    data_type : numpy.dtype-like | list[numpy.dtype-like], optional
        If set, find all channels with these data types.
    min_rate : int, optional
        The minimum sampling rate for channels.
    max_rate : int, optional
        The maximum sampling rate for channels.
    publisher : str | list[str], optional
        If set, find all channels associated with these publishers.
    time : float, optional
        GPS time in seconds indicating when the metadata query is valid.
        If None, routes to the live backend (current state).
    replay_id : str, optional
        Server-registered replay identifier.

    Returns
    -------
    int
        The number of channels matching query.

    """
    data_type = _parse_data_types(data_type)
    if min_rate is None:
        min_rate = constants.MIN_SAMPLE_RATE
    if max_rate is None:
        max_rate = constants.MAX_SAMPLE_RATE
    if publisher is None:
        publisher = []
    elif isinstance(publisher, str):
        publisher = [publisher]

    time_ns = time_as_ns(time) if time is not None else None
    descriptor = create_descriptor(
        RequestType.Count,
        pattern=pattern,
        data_type=data_type,
        min_rate=min_rate,
        max_rate=max_rate,
        publisher=publisher,
        time=time_ns,
        replay_id=replay_id,
        validator=self._validator,
    )
    count = 0
    with connect(self.initial_url, middleware=self._middleware) as client:
        flight_info = get_flight_info(client, descriptor)
        with MultiFlightReader(
            flight_info.endpoints, client, middleware=self._middleware
        ) as stream:
            for data in stream.unpack():
                count += data["count"]
    return count

describe

describe(channels, time=None, replay_id=None)

Get channel metadata for channels requested

Parameters:

Name Type Description Default
channels list[str]

List of channels to request.

required
time float

GPS time in seconds indicating when the metadata query is valid. If None, routes to the live backend (current state).

None
replay_id str

Server-registered replay identifier.

None

Returns:

Type Description
dict[str, Channel]

Mapping of channel names to channel metadata.

Source code in arrakis/client.py
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
def describe(
    self,
    channels: list[str],
    time: float | None = None,
    replay_id: str | None = None,
) -> dict[str, Channel]:
    """Get channel metadata for channels requested

    Parameters
    ----------
    channels : list[str]
        List of channels to request.
    time : float, optional
        GPS time in seconds indicating when the metadata query is valid.
        If None, routes to the live backend (current state).
    replay_id : str, optional
        Server-registered replay identifier.

    Returns
    -------
    dict[str, Channel]
        Mapping of channel names to channel metadata.

    """
    time_ns = time_as_ns(time) if time is not None else None
    with connect(self.initial_url, middleware=self._middleware) as client:
        return self._describe(
            client,
            channels,
            time_ns=time_ns,
            replay_id=replay_id,
        )

domains

domains()

Get the list of domains available on the server.

Returns:

Type Description
list[str]

Sorted list of domain names.

Source code in arrakis/client.py
588
589
590
591
592
593
594
595
596
597
598
599
600
601
def domains(self) -> list[str]:
    """Get the list of domains available on the server.

    Returns
    -------
    list[str]
        Sorted list of domain names.

    """
    result = self.scope_map()
    domain_set = set()
    for endpoint_info in result["endpoints"].values():
        domain_set.update(endpoint_info["scopes"].keys())
    return sorted(domain_set)

fetch

fetch(channels, start, end, replay_start=None, replay_end=None, replay_id=None, on_gap=ONGAP_DEFAULT)

Fetch timeseries data

Parameters:

Name Type Description Default
channels list[str]

List of channels to request.

required
start float

GPS start time, in seconds.

required
end float

GPS end time, in seconds.

required
replay_start float

GPS start time of the replay window, in seconds.

None
replay_end float

GPS end time of the replay window, in seconds.

None
replay_id str

Server-registered replay identifier. Mutually exclusive with replay_start/replay_end.

None
on_gap str

Policy for spans with no data: 'fill' (default) returns masked arrays over the gaps (note the memory cost over large holes), 'raise' raises arrakis.mux.GapError. 'skip' is not supported: fetch returns one contiguous block.

ONGAP_DEFAULT

Returns:

Type Description
SeriesBlock

Dictionary-like object containing all requested channel data.

Source code in arrakis/client.py
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
def fetch(
    self,
    channels: list[str],
    start: float,
    end: float,
    replay_start: float | None = None,
    replay_end: float | None = None,
    replay_id: str | None = None,
    on_gap: str = ONGAP_DEFAULT,
) -> SeriesBlock:
    """Fetch timeseries data

    Parameters
    ----------
    channels : list[str]
        List of channels to request.
    start : float
        GPS start time, in seconds.
    end : float
        GPS end time, in seconds.
    replay_start : float, optional
        GPS start time of the replay window, in seconds.
    replay_end : float, optional
        GPS end time of the replay window, in seconds.
    replay_id : str, optional
        Server-registered replay identifier. Mutually exclusive with
        replay_start/replay_end.
    on_gap : str, optional
        Policy for spans with no data: 'fill' (default) returns
        masked arrays over the gaps (note the memory cost over
        large holes), 'raise' raises ``arrakis.mux.GapError``.
        'skip' is not supported: fetch returns one contiguous
        block.

    Returns
    -------
    SeriesBlock
        Dictionary-like object containing all requested channel data.

    """
    if on_gap.upper() == "SKIP":
        msg = "fetch returns one contiguous block; on_gap='skip' is not supported"
        raise ValueError(msg)
    return concatenate_blocks(
        *self.stream(
            channels,
            start,
            end,
            replay_start=replay_start,
            replay_end=replay_end,
            replay_id=replay_id,
            on_gap=on_gap,
        )
    )

find

find(pattern=constants.DEFAULT_MATCH, data_type=None, min_rate=constants.MIN_SAMPLE_RATE, max_rate=constants.MAX_SAMPLE_RATE, publisher=None, time=None, replay_id=None)

Find channels matching a set of conditions

Parameters:

Name Type Description Default
pattern str

Channel pattern to match channels with, using regular expressions.

DEFAULT_MATCH
data_type dtype - like | list[dtype - like]

If set, find all channels with these data types.

None
min_rate int

Minimum sampling rate for channels.

MIN_SAMPLE_RATE
max_rate int

Maximum sampling rate for channels.

MAX_SAMPLE_RATE
publisher str | list[str]

If set, find all channels associated with these publishers.

None
time float

GPS time in seconds indicating when the metadata query is valid. If None, routes to the live backend (current state).

None
replay_id str

Server-registered replay identifier.

None

Yields:

Type Description
Channel

Channel objects for all channels matching query.

Source code in arrakis/client.py
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
def find(
    self,
    pattern: str = constants.DEFAULT_MATCH,
    data_type: DataTypeLike | None = None,
    min_rate: int | None = constants.MIN_SAMPLE_RATE,
    max_rate: int | None = constants.MAX_SAMPLE_RATE,
    publisher: str | list[str] | None = None,
    time: float | None = None,
    replay_id: str | None = None,
) -> Generator[Channel, None, None]:
    """Find channels matching a set of conditions

    Parameters
    ----------
    pattern : str, optional
        Channel pattern to match channels with, using regular expressions.
    data_type : numpy.dtype-like | list[numpy.dtype-like], optional
        If set, find all channels with these data types.
    min_rate : int, optional
        Minimum sampling rate for channels.
    max_rate : int, optional
        Maximum sampling rate for channels.
    publisher : str | list[str], optional
        If set, find all channels associated with these publishers.
    time : float, optional
        GPS time in seconds indicating when the metadata query is valid.
        If None, routes to the live backend (current state).
    replay_id : str, optional
        Server-registered replay identifier.

    Yields
    -------
    Channel
        Channel objects for all channels matching query.

    """
    data_type = _parse_data_types(data_type)
    if min_rate is None:
        min_rate = constants.MIN_SAMPLE_RATE
    if max_rate is None:
        max_rate = constants.MAX_SAMPLE_RATE
    if publisher is None:
        publisher = []
    elif isinstance(publisher, str):
        publisher = [publisher]

    time_ns = time_as_ns(time) if time is not None else None
    descriptor = create_descriptor(
        RequestType.Find,
        pattern=pattern,
        data_type=data_type,
        min_rate=min_rate,
        max_rate=max_rate,
        publisher=publisher,
        time=time_ns,
        replay_id=replay_id,
        validator=self._validator,
    )
    with connect(self.initial_url, middleware=self._middleware) as client:
        yield from self._stream_channel_metadata(client, descriptor)

replays

replays()

Get available replay windows with start/end times.

Returns:

Type Description
dict

Mapping of replay IDs to their start/end GPS times in seconds.

Source code in arrakis/client.py
567
568
569
570
571
572
573
574
575
576
def replays(self) -> dict:
    """Get available replay windows with start/end times.

    Returns
    -------
    dict
        Mapping of replay IDs to their start/end GPS times in seconds.

    """
    return self._do_action("replays")

scope_map

scope_map()

Get the mapping of endpoints to their scopes.

Returns:

Type Description
dict

Mapping of endpoint information including scopes per domain.

Source code in arrakis/client.py
556
557
558
559
560
561
562
563
564
565
def scope_map(self) -> dict:
    """Get the mapping of endpoints to their scopes.

    Returns
    -------
    dict
        Mapping of endpoint information including scopes per domain.

    """
    return self._do_action("scope-map")

server_info

server_info()

Get server version and capability metadata.

Returns:

Type Description
dict

Server metadata including version info, backend type, capabilities, and domains.

Source code in arrakis/client.py
544
545
546
547
548
549
550
551
552
553
554
def server_info(self) -> dict:
    """Get server version and capability metadata.

    Returns
    -------
    dict
        Server metadata including version info, backend type,
        capabilities, and domains.

    """
    return self._do_action("server-info")

stream

stream(channels, start=None, end=None, kafka_url=None, replay_start=None, replay_end=None, replay_id=None, on_gap=ONGAP_DEFAULT, stride=None, *, fast_forward=True)

Stream live or offline timeseries data

Parameters:

Name Type Description Default
channels list[str]

List of channels to request.

required
start float

GPS start time, in seconds.

None
end float

GPS end time, in seconds.

None
replay_start float

GPS start time of the replay window, in seconds.

None
replay_end float

GPS end time of the replay window, in seconds.

None
replay_id str

Server-registered replay identifier. Mutually exclusive with replay_start/replay_end.

None
on_gap str

Policy for spans with no data: 'fill' (default) yields masked gap blocks so the stream is continuous, 'skip' yields nothing for such spans (sparse blocks), 'raise' raises arrakis.mux.GapError.

ONGAP_DEFAULT
stride float

Duration of yielded blocks, in seconds. Must be an integer multiple of the stream's native stride (for multiple channels, the least common multiple of the channels' strides), which is also the default. Blocks are aggregated client-side: live blocks are yielded only once their full span has been received, and align to absolute GPS multiples of the stride; bounded requests align to start, with a shorter final block if the requested span is not a multiple of the stride.

None
fast_forward bool

When True (default), a live stream read directly from Kafka that has fallen further behind real time than its latency budget allows -- for long enough that it clearly cannot catch up -- skips ahead to current data, trading the unusable backlog for a gap. The skip and recovery are logged. Ignored for bounded requests (reading past data is their job) and for server-relayed streams (the server manages those itself).

True

Yields:

Type Description
SeriesBlock

Dictionary-like object containing all requested channel data.

Setting neither start nor end begins a live stream starting
from now.
Source code in arrakis/client.py
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
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
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
def stream(
    self,
    channels: list[str],
    start: float | None = None,
    end: float | None = None,
    kafka_url: str | None = None,
    replay_start: float | None = None,
    replay_end: float | None = None,
    replay_id: str | None = None,
    on_gap: str = ONGAP_DEFAULT,
    stride: float | None = None,
    *,
    fast_forward: bool = True,
) -> Generator[SeriesBlock, None, None]:
    """Stream live or offline timeseries data

    Parameters
    ----------
    channels : list[str]
        List of channels to request.
    start : float, optional
        GPS start time, in seconds.
    end : float, optional
        GPS end time, in seconds.
    replay_start : float, optional
        GPS start time of the replay window, in seconds.
    replay_end : float, optional
        GPS end time of the replay window, in seconds.
    replay_id : str, optional
        Server-registered replay identifier. Mutually exclusive with
        replay_start/replay_end.
    on_gap : str, optional
        Policy for spans with no data: 'fill' (default) yields
        masked gap blocks so the stream is continuous, 'skip'
        yields nothing for such spans (sparse blocks), 'raise'
        raises ``arrakis.mux.GapError``.
    stride : float, optional
        Duration of yielded blocks, in seconds.  Must be an
        integer multiple of the stream's native stride (for
        multiple channels, the least common multiple of the
        channels' strides), which is also the default.  Blocks
        are aggregated client-side: live blocks are yielded only
        once their full span has been received, and align to
        absolute GPS multiples of the stride; bounded requests
        align to ``start``, with a shorter final block if the
        requested span is not a multiple of the stride.
    fast_forward : bool, optional
        When True (default), a live stream read directly from
        Kafka that has fallen further behind real time than its
        latency budget allows -- for long enough that it clearly
        cannot catch up -- skips ahead to current data, trading
        the unusable backlog for a gap.  The skip and recovery
        are logged.  Ignored for bounded requests (reading past
        data is their job) and for server-relayed streams (the
        server manages those itself).

    Yields
    ------
    SeriesBlock
        Dictionary-like object containing all requested channel data.

    Setting neither start nor end begins a live stream starting
    from now.

    """
    start_ns = time_as_ns(start) if start is not None else None
    end_ns = time_as_ns(end) if end is not None else None
    stride_ns = time_as_ns(stride) if stride is not None else None
    replay_start_ns = time_as_ns(replay_start) if replay_start is not None else None
    replay_end_ns = time_as_ns(replay_end) if replay_end is not None else None
    metadata: dict[str, Channel] = {}
    reader: StreamReader

    # for opportunistic replay, translate start time into the
    # replay window so describe routes to the archival backend
    describe_time: int | None
    if replay_start_ns is not None and replay_end_ns is not None and not replay_id:
        describe_time = translate_time(start_ns, replay_start_ns, replay_end_ns)
    else:
        describe_time = start_ns

    with connect(self.initial_url, middleware=self._middleware) as client:
        # initial describe request to get the full metadata needed
        # to construct the appropriate muxer — use start time so the
        # info server routes to the correct backend for the request
        metadata = self._describe(
            client,
            channels,
            time_ns=describe_time,
            replay_id=replay_id,
        )
        for channel in metadata.values():
            assert channel.stride, "Channels do not include stride info."

        # get the flight info for streams
        descriptor = create_descriptor(
            RequestType.Stream,
            channels=channels,
            start=start_ns,
            end=end_ns,
            replay_start=replay_start_ns,
            replay_end=replay_end_ns,
            replay_id=replay_id,
            validator=self._validator,
        )
        flight_info = get_flight_info(client, descriptor)

    logger.info(
        "streaming %d channel(s), %s",
        len(channels),
        _describe_span(start, end, replay_id, replay_start, replay_end),
    )

    reader = _construct_stream_reader(
        flight_info.endpoints,
        metadata,
        start_ns,
        kafka_url=kafka_url,
        initial_url=self.initial_url,
        middleware=self._middleware,
        # live streams only: a reader over past data is expected
        # to be behind
        fast_forward=fast_forward and start_ns is None,
    )

    # streams served through flight endpoints are relays of a
    # server-side muxer and need extra timeout headroom; direct
    # kafka reads consume the source and keep per-queue defaults
    mux_timeout = 0
    if kafka_url is None and any(
        not _endpoint_is_kafka(ep) for ep in flight_info.endpoints
    ):
        mux_timeout = _relay_mux_timeout(metadata)

    mux: BlockMuxStream = BlockMuxStream(
        reader, start=start_ns, timeout=mux_timeout, on_gap=on_gap, stride=stride_ns
    )

    for block in mux.stream(end_ns):
        yield block

parse_arrakis_url

parse_arrakis_url(url=None)

Parse an Arrakis URL, filling in with default values

If url is None or unspecified, the ARRAKIS_SERVER environment variable will be used. If the port is not specified, 31206 will be used.

Source code in arrakis/client.py
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
def parse_arrakis_url(url: str | None = None) -> ParseResult:
    """Parse an Arrakis URL, filling in with default values

    If url is None or unspecified, the ARRAKIS_SERVER environment
    variable will be used.  If the port is not specified, 31206 will
    be used.

    """
    if url is None:
        url = os.getenv("ARRAKIS_SERVER", constants.DEFAULT_ARRAKIS_SERVER)
    assert url is not None, "ARRAKIS_SERVER not specified."
    if "://" not in url:
        url = "grpc://" + url
    parsed = urlparse(url)
    if ":" not in parsed.netloc:
        parsed = parsed._replace(netloc=parsed.netloc + ":31206")
    if parsed.netloc[0] == ":":
        parsed = parsed._replace(netloc="localhost" + parsed.netloc)
    return parsed