Skip to content

fujilib.streaming

One Sample per poll carries a whole Frame. A poll source (PollSourceAdapter for one analyzer) is what record() polls; the recorder yields a Recording of per-tick batches with a live AcquisitionSummary (design §7.6). The guide is Recording.

fujilib.streaming.sample

Timed sample: one poll of one analyzer, with its timing provenance (unified API §C).

One :class:Sample per analyzer per tick carries the whole :class:Frame (design §7.6, §13.1 #13), so analyzer-level status travels once with every channel. The unified timestamp contract:

  • :attr:Sample.t_mono_ns — the monotonic acquisition time, the join key. It is the request/reply midpoint of the block that holds every concentration.
  • :attr:Sample.t_utc — the same instant on the wall clock (UTC, tz-aware).
  • :attr:Sample.t_midpoint_mono_ns — an integration-window midpoint. Always None: a configured averaging period does not reveal the actual integration window, so response-time and averaging settings are reported in the analyzer metadata instead of shifting timestamps.

requested_at, received_at and latency_s are those of the same concentration block, so t_utc is their midpoint. The status block's timing is in Frame.status_timing.

A failed poll is still a sample: frame is None, error is set, and the timing of the attempt is kept, so gaps are recorded rather than dropped.

:attr:Sample.channels names the channels a row of the sample has columns for. A recorder sets it to the channels established when the recording started, on every sample it makes, successful or not, so sample_to_row(sample) gives every row of a recording the same keys.

Sample dataclass

Sample(
    device,
    address,
    frame,
    protocol,
    t_mono_ns,
    t_utc,
    requested_at,
    received_at,
    latency_s,
    t_midpoint_mono_ns=None,
    metadata=_empty_metadata(),
    error=None,
    channels=(),
)

One poll of one analyzer, with full timing provenance.

Attributes:

Name Type Description
device str

The name that follows the data into sinks.

address int

The analyzer's station number.

frame Frame | None

Every established channel and the analyzer status, or None when the poll failed.

protocol ProtocolKind

The wire protocol, kept for error rows too.

t_mono_ns int

Monotonic acquisition time: the concentration block's request/reply midpoint, in nanoseconds.

t_utc datetime

The same instant on the wall clock.

requested_at datetime

Wall clock just before the concentration block was requested.

received_at datetime

Wall clock just after its reply (or failure) was seen.

latency_s float

received_at - requested_at in seconds.

t_midpoint_mono_ns int | None

Integration-window midpoint; always None.

metadata Mapping[str, str]

Free-form annotations. Not written to rows.

error FujiError | None

The error of a failed poll, or None.

channels tuple[ChannelId, ...]

The channels a row of this sample has columns for, in order. A recorder fixes them when it starts (design §13.1 #34).

from_error classmethod

from_error(
    error,
    *,
    device,
    address,
    protocol,
    timing,
    metadata=None,
    channels=(),
)

Build the sample of a failed poll, timed by the failed attempt.

Pass the recording's channels so the sample's row has the same columns as a successful one.

Source code in src/fujilib/streaming/sample.py
@classmethod
def from_error(
    cls,
    error: FujiError,
    *,
    device: str,
    address: int,
    protocol: ProtocolKind,
    timing: TransferTiming,
    metadata: Mapping[str, str] | None = None,
    channels: Iterable[ChannelId] = (),
) -> Sample:
    """Build the sample of a failed poll, timed by the failed attempt.

    Pass the recording's ``channels`` so the sample's row has the same
    columns as a successful one.
    """
    return cls(
        device=device,
        address=address,
        frame=None,
        protocol=protocol,
        t_mono_ns=timing.midpoint_mono_ns,
        t_utc=timing.midpoint_utc,
        requested_at=timing.requested_at,
        received_at=timing.received_at,
        latency_s=timing.latency_s,
        metadata=MappingProxyType(dict(metadata or {})),
        error=error,
        channels=tuple(channels),
    )

from_frame classmethod

from_frame(
    frame, *, device, address, metadata=None, channels=None
)

Build the sample of a successful poll, timed by its concentration block.

channels defaults to the frame's own.

Source code in src/fujilib/streaming/sample.py
@classmethod
def from_frame(
    cls,
    frame: Frame,
    *,
    device: str,
    address: int,
    metadata: Mapping[str, str] | None = None,
    channels: Iterable[ChannelId] | None = None,
) -> Sample:
    """Build the sample of a successful poll, timed by its concentration block.

    ``channels`` defaults to the frame's own.
    """
    timing = frame.readings_timing
    return cls(
        device=device,
        address=address,
        frame=frame,
        protocol=frame.protocol,
        t_mono_ns=timing.midpoint_mono_ns,
        t_utc=timing.midpoint_utc,
        requested_at=timing.requested_at,
        received_at=timing.received_at,
        latency_s=timing.latency_s,
        metadata=MappingProxyType(dict(metadata or {})),
        channels=frame.channels if channels is None else tuple(channels),
    )

fujilib.streaming.poll_source

Poll sources: DeviceResult, PollSource and PollSourceAdapter (unified API §E).

A recorder polls a poll source: an object whose poll(names) returns one :class:DeviceResult per analyzer, keyed by name, holding either a :class:~fujilib.devices.models.Frame or the error of a failed poll. :class:PollSourceAdapter makes one :class:~fujilib.devices.analyzer.Analyzer such a source, under a name that follows its data into rows.

A source also describes each analyzer with a :class:SourceLayout: its station, protocol and established channels. The recorder reads the layouts once, when it starts, so every sample of a recording, successful or not, has the same row columns (design §7.6, §13.1 #34).

DeviceResult dataclass

DeviceResult(value, error)

One device's outcome: a value or an error, never both.

Attributes:

Name Type Description
value T | None

The result, or None when the call failed.

error FujiError | None

The error, or None when the call succeeded.

ok property

ok

Whether the device produced a value.

failure classmethod

failure(error)

A failed result holding error.

Source code in src/fujilib/streaming/poll_source.py
@classmethod
def failure(cls, error: FujiError) -> Self:
    """A failed result holding ``error``."""
    return cls(value=None, error=error)

success classmethod

success(value)

A successful result holding value.

Source code in src/fujilib/streaming/poll_source.py
@classmethod
def success[V](cls, value: V) -> DeviceResult[V]:
    """A successful result holding ``value``."""
    del cls
    return DeviceResult(value=value, error=None)

PollSource

Bases: Protocol

What :func:~fujilib.streaming.recorder.record polls (design §7.6).

layout

layout(names=None)

Describe the named analyzers (all when None), with no I/O.

Source code in src/fujilib/streaming/poll_source.py
def layout(self, names: Sequence[str] | None = None) -> Mapping[str, SourceLayout]:
    """Describe the named analyzers (all when ``None``), with no I/O."""
    ...

poll async

poll(names=None)

Poll the named analyzers (all when None) once; a failure is a failed result.

Source code in src/fujilib/streaming/poll_source.py
async def poll(self, names: Sequence[str] | None = None) -> Mapping[str, DeviceResult[Frame]]:
    """Poll the named analyzers (all when ``None``) once; a failure is a failed result."""
    ...

reconnect async

reconnect(name)

Open analyzer name again after a connection failure.

Raises:

Type Description
FujiError

it could not be opened, or is not the analyzer it was.

Source code in src/fujilib/streaming/poll_source.py
async def reconnect(self, name: str) -> None:
    """Open analyzer ``name`` again after a connection failure.

    Raises:
        FujiError: it could not be opened, or is not the analyzer it was.
    """
    ...

PollSourceAdapter

PollSourceAdapter(name, device)

One :class:~fujilib.devices.analyzer.Analyzer as a poll source.

Example::

source = PollSourceAdapter("zpa", analyzer)
results = await source.poll()  # {"zpa": DeviceResult(frame, None)}

The parameter is device, as in every sibling library (unified API §E).

Publish device's polls under name.

Source code in src/fujilib/streaming/poll_source.py
def __init__(self, name: str, device: Analyzer) -> None:
    """Publish ``device``'s polls under ``name``."""
    self._name = name
    self._device = device

device property

device

The wrapped analyzer.

name property

name

The name the analyzer's data is published under.

layout

layout(names=None)

The analyzer's station, protocol and established channels, under its name.

An empty mapping when names leaves this analyzer out.

Source code in src/fujilib/streaming/poll_source.py
def layout(self, names: Sequence[str] | None = None) -> Mapping[str, SourceLayout]:
    """The analyzer's station, protocol and established channels, under its name.

    An empty mapping when ``names`` leaves this analyzer out.
    """
    if not self._selected(names):
        return MappingProxyType({})
    device = self._device
    layout = SourceLayout(
        address=device.address,
        protocol=device.protocol,
        channels=tuple(c.channel for c in device.channels),
        reopenable=device.session.reopenable,
    )
    return MappingProxyType({self._name: layout})

poll async

poll(names=None)

Poll the analyzer once: two transactions.

Returns {name: result}, or an empty mapping when names leaves this analyzer out. A failed poll is a failed result, not an exception.

Source code in src/fujilib/streaming/poll_source.py
async def poll(self, names: Sequence[str] | None = None) -> Mapping[str, DeviceResult[Frame]]:
    """Poll the analyzer once: two transactions.

    Returns ``{name: result}``, or an empty mapping when ``names`` leaves
    this analyzer out. A failed poll is a failed result, not an exception.
    """
    if not self._selected(names):
        return MappingProxyType({})
    result: DeviceResult[Frame]
    try:
        result = DeviceResult.success(await self._device.poll())
    except FujiError as exc:
        result = DeviceResult(value=None, error=exc)
    return MappingProxyType({self._name: result})

reconnect async

reconnect(name)

Reopen the analyzer (:meth:Analyzer.reopen <fujilib.devices.analyzer.Analyzer.reopen>).

Raises:

Type Description
FujiValidationError

name is not this source's name.

FujiError

the analyzer could not be reopened.

Source code in src/fujilib/streaming/poll_source.py
async def reconnect(self, name: str) -> None:
    """Reopen the analyzer (:meth:`Analyzer.reopen <fujilib.devices.analyzer.Analyzer.reopen>`).

    Raises:
        FujiValidationError: ``name`` is not this source's name.
        FujiError: the analyzer could not be reopened.
    """
    if name != self._name:
        msg = f"this source publishes {self._name!r}, not {name!r}"
        raise FujiValidationError(msg)
    _ = await self._device.reopen()

SourceLayout dataclass

SourceLayout(address, protocol, channels, reopenable=False)

What a recorder needs to know about one analyzer of a source before it polls.

Attributes:

Name Type Description
address int

The station number, 1-31.

protocol ProtocolKind

The wire protocol.

channels tuple[ChannelId, ...]

The established channels, in order: the row columns.

reopenable bool

Whether the source can open the analyzer again after a connection failure (:meth:PollSource.reconnect).

fujilib.streaming.recorder

The recorder: poll a source at a fixed rate into a bounded stream (design §7.6).

:func:record is an async context manager. It yields a :class:Recording whose stream receives one batch per tick: a mapping from each analyzer's name to its :class:~fujilib.streaming.sample.Sample::

async with record(PollSourceAdapter("zpa", anz), rate_hz=1.0, duration=60) as rec:
    async for batch in rec.stream:
        sample = batch["zpa"]

Rows keep their columns. The source's layouts are read once, when the recording starts, and every sample carries the channels established then, so sample_to_row(sample) gives every row the same keys, a failed poll included. A channel established later is left out of this recording's rows (design §13.1 #34).

Schedule. Ticks follow an absolute schedule: tick k is due k / rate_hz seconds after the first, which runs at once. A poll that overruns by one or more whole slots skips them and counts them in samples_late; the recorder never catches up in a burst. max_drift_ms is the latest a poll started after its slot. A recording with a duration makes every tick due before the duration has passed, ceil(duration * rate_hz) of them (at least one), and its stream ends after the last.

Failures. A failed poll is a sample with frame=None and error set, timed around the poll, so gaps are recorded; the recording goes on. A connection failure (:class:~fujilib.errors.FujiConnectionError) ends it: the tick's batch is still delivered, the stream ends, and leaving the async with block raises the error. With a :class:ReconnectPolicy the recording goes on instead: every tick of the outage is an error sample, and the analyzer is reopened on the policy's schedule.

Overflow. The stream holds buffer_size batches. When it is full, :class:OverflowPolicy decides: wait (BLOCK, the consumer sets the pace and ticks missed meanwhile are late), or drop a whole batch, the new one (DROP_NEWEST) or the oldest (DROP_OLDEST), counted in samples_dropped. A batch is never split. Batches still waiting in the stream when the recording stops were never read, and are counted as dropped too, so samples_emitted is then exactly what the consumer received.

For a recording that runs to its end, samples_emitted + samples_dropped + samples_late == target_total_samples.

Batch

Batch = Mapping[str, Sample]

One tick of a recording: each analyzer's sample, by name.

AcquisitionSummary dataclass

AcquisitionSummary(
    started_at,
    finished_at=None,
    samples_emitted=0,
    samples_late=0,
    max_drift_ms=0.0,
    target_total_samples=None,
    samples_dropped=0,
    error_samples=0,
    disconnects=0,
    reconnects=0,
)

The counters of one recording (unified API §I, §M).

Mutable: the recorder updates it while it runs, and sets finished_at when it stops, however it stops. Every count is of polls (batches), not of rows (design §7.6).

disconnects class-attribute instance-attribute

disconnects = 0

Connection failures: each outage a ReconnectPolicy rides out, and the one that ends a recording without one.

error_samples class-attribute instance-attribute

error_samples = 0

Samples of failed polls (frame=None), counted as they are polled.

max_drift_ms class-attribute instance-attribute

max_drift_ms = 0.0

The latest a poll started after its slot, in milliseconds.

reconnects class-attribute instance-attribute

reconnects = 0

Outages that ended with the analyzer reopened.

samples_dropped class-attribute instance-attribute

samples_dropped = 0

Batches the overflow policy discarded, and batches still unread when the recording stopped.

samples_emitted class-attribute instance-attribute

samples_emitted = 0

Batches put on the stream for the consumer and not dropped since; once the recording has stopped, exactly the batches the consumer received.

samples_late class-attribute instance-attribute

samples_late = 0

Ticks skipped because a poll overran their slots.

target_total_samples class-attribute instance-attribute

target_total_samples = None

Ticks a recording with a duration is due to make; None without one.

OverflowPolicy

Bases: StrEnum

What the recorder does when the stream's buffer is full.

BLOCK class-attribute instance-attribute

BLOCK = 'block'

Wait for room. The consumer sets the pace; ticks missed meanwhile are late.

DROP_NEWEST class-attribute instance-attribute

DROP_NEWEST = 'drop_newest'

Discard the new batch.

DROP_OLDEST class-attribute instance-attribute

DROP_OLDEST = 'drop_oldest'

Discard the oldest waiting batch to make room for the new one.

ReconnectPolicy dataclass

ReconnectPolicy(
    backoff_s=DEFAULT_BACKOFF_S, max_attempts=None
)

Reopen an analyzer after a connection failure instead of ending the recording.

The recording keeps its schedule through an outage: every tick is an error sample until the analyzer is back. Attempts are made at ticks, the first backoff_s[0] seconds after the failure, the next backoff_s[1] after that, and so on; the last wait repeats. An attempt takes as long as opening the port and identifying the analyzer, and delays the tick it is made in.

An attempt that finds another analyzer, or any failure other than a connection, timeout or framing error, ends the recording.

Attributes:

Name Type Description
backoff_s tuple[float, ...]

Seconds to wait before each attempt.

max_attempts int | None

Attempts per outage before the recording ends with the last error; None for no limit.

__post_init__

__post_init__()

Check the schedule.

Raises:

Type Description
FujiValidationError

an empty or negative schedule, or fewer than one attempt.

Source code in src/fujilib/streaming/recorder.py
def __post_init__(self) -> None:
    """Check the schedule.

    Raises:
        FujiValidationError: an empty or negative schedule, or fewer than one attempt.
    """
    if not self.backoff_s or not all(_is_seconds(s) and s >= 0 for s in self.backoff_s):
        msg = f"backoff_s must be one or more finite, non-negative seconds: {self.backoff_s!r}"
        raise FujiValidationError(msg)
    attempts = self.max_attempts
    if attempts is not None and (not _is_int(attempts) or attempts < 1):
        msg = f"max_attempts must be None or at least 1, got {attempts!r}"
        raise FujiValidationError(msg)

delay

delay(attempt)

The wait before attempt attempt (1 for the first) of an outage.

Source code in src/fujilib/streaming/recorder.py
def delay(self, attempt: int) -> float:
    """The wait before attempt ``attempt`` (1 for the first) of an outage."""
    return self.backoff_s[min(attempt, len(self.backoff_s)) - 1]

Recording dataclass

Recording(stream, summary, rate_hz)

A running recording (unified API §I).

Attributes:

Name Type Description
stream MemoryObjectReceiveStream[T]

The batches, one per tick; iterate it, or the recording itself.

summary AcquisitionSummary

The live counters.

rate_hz float

The requested rate.

__aiter__

__aiter__()

Iterate the batches of :attr:stream.

Source code in src/fujilib/streaming/recorder.py
def __aiter__(self) -> AsyncIterator[T]:
    """Iterate the batches of :attr:`stream`."""
    return self.stream.__aiter__()

record async

record(
    source,
    *,
    rate_hz,
    duration=None,
    names=None,
    overflow=OverflowPolicy.BLOCK,
    buffer_size=64,
    reconnect=None,
)

Poll source at rate_hz for the async with block (see the module docstring).

Parameters:

Name Type Description Default
source PollSource

The analyzers, e.g. a :class:~fujilib.streaming.poll_source.PollSourceAdapter. Each must have established channels: identify it, or assert its gases with channel_map.

required
rate_hz float

Ticks per second. One analyzer's poll takes about 0.12 s, so about 7-8 Hz is the practical ceiling for one station (design §2.4).

required
duration float | None

Seconds to record; None records until the block exits.

None
names Sequence[str] | None

The analyzers to record; all of the source's when None.

None
overflow OverflowPolicy

What to do when buffer_size batches are waiting.

BLOCK
buffer_size int

Batches the stream holds.

64
reconnect ReconnectPolicy | None

Reopen an analyzer after a connection failure instead of ending the recording; each analyzer must be reopenable.

None

Raises:

Type Description
FujiValidationError

an argument is invalid, a name is unknown, or an analyzer has no established channels (or cannot be reopened, with reconnect); nothing was polled.

FujiConnectionError

raised on leaving the block when a connection failure ended the recording.

Source code in src/fujilib/streaming/recorder.py
@asynccontextmanager
async def record(
    source: PollSource,
    *,
    rate_hz: float,
    duration: float | None = None,
    names: Sequence[str] | None = None,
    overflow: OverflowPolicy = OverflowPolicy.BLOCK,
    buffer_size: int = 64,
    reconnect: ReconnectPolicy | None = None,
) -> AsyncGenerator[Recording[Batch]]:
    """Poll ``source`` at ``rate_hz`` for the ``async with`` block (see the module docstring).

    Args:
        source: The analyzers, e.g. a
            :class:`~fujilib.streaming.poll_source.PollSourceAdapter`. Each must
            have established channels: identify it, or assert its gases with
            ``channel_map``.
        rate_hz: Ticks per second. One analyzer's poll takes about 0.12 s, so
            about 7-8 Hz is the practical ceiling for one station (design §2.4).
        duration: Seconds to record; ``None`` records until the block exits.
        names: The analyzers to record; all of the source's when ``None``.
        overflow: What to do when ``buffer_size`` batches are waiting.
        buffer_size: Batches the stream holds.
        reconnect: Reopen an analyzer after a connection failure instead of
            ending the recording; each analyzer must be reopenable.

    Raises:
        FujiValidationError: an argument is invalid, a name is unknown, or an
            analyzer has no established channels (or cannot be reopened, with
            ``reconnect``); nothing was polled.
        FujiConnectionError: raised on leaving the block when a connection
            failure ended the recording.
    """
    async with _record(
        source,
        rate_hz=rate_hz,
        duration=duration,
        names=names,
        overflow=overflow,
        buffer_size=buffer_size,
        reconnect=reconnect,
        clock=_ANYIO_CLOCK,
    ) as recording:
        yield recording

fujilib.units

Pint-compatible unit strings — :func:to_pint (unified API §K).

Every sibling library exposes the same free function to_pint(unit) -> str | None in <lib>.units, so consumers can resolve units without knowing which device produced a row. pint is not a dependency: this returns plain strings, and the consumer parses them.

The strings are chosen so capa's unit registry parses and canonicalizes them. vol% is a volume fraction in percent, so it maps to "percent"; ppm maps to "ppm". The two are dimensionally compatible and differ by 10⁴, so a consumer must compare canonical unit strings, not only dimensions. to_pint never converts values.

Unit

Bases: StrEnum

A concentration unit. The value is the canonical display string.

to_pint

to_pint(unit)

Return a pint-compatible unit string for unit, or None.

Accepts a :class:Unit, a unit string such as "vol%" or "ppm", or None. Unknown strings, :attr:Unit.UNKNOWN and None return None; this never raises.

Example::

>>> to_pint(Unit.VOL_PERCENT)
'percent'
>>> to_pint("mg/m3")
'mg/m**3'
>>> to_pint("furlongs")
Source code in src/fujilib/units.py
def to_pint(unit: Unit | str | None) -> str | None:
    """Return a pint-compatible unit string for ``unit``, or ``None``.

    Accepts a :class:`Unit`, a unit string such as ``"vol%"`` or ``"ppm"``, or
    ``None``. Unknown strings, :attr:`Unit.UNKNOWN` and ``None`` return
    ``None``; this never raises.

    Example::

        >>> to_pint(Unit.VOL_PERCENT)
        'percent'
        >>> to_pint("mg/m3")
        'mg/m**3'
        >>> to_pint("furlongs")
    """
    if unit is None:
        return None
    return _UNIT_TO_PINT[coerce_unit(unit)]