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. AlwaysNone: 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 |
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
|
|
t_midpoint_mono_ns |
int | None
|
Integration-window midpoint; always |
metadata |
Mapping[str, str]
|
Free-form annotations. Not written to rows. |
error |
FujiError | None
|
The error of a failed poll, or |
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
¶
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
from_frame
classmethod
¶
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
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
¶
PollSource ¶
Bases: Protocol
What :func:~fujilib.streaming.recorder.record polls (design §7.6).
PollSourceAdapter ¶
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
layout ¶
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
poll
async
¶
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
reconnect
async
¶
Reopen the analyzer (:meth:Analyzer.reopen <fujilib.devices.analyzer.Analyzer.reopen>).
Raises:
| Type | Description |
|---|---|
FujiValidationError
|
|
FujiError
|
the analyzer could not be reopened. |
Source code in src/fujilib/streaming/poll_source.py
SourceLayout
dataclass
¶
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: |
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.
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
¶
Connection failures: each outage a ReconnectPolicy rides out, and the one
that ends a recording without one.
error_samples
class-attribute
instance-attribute
¶
Samples of failed polls (frame=None), counted as they are polled.
max_drift_ms
class-attribute
instance-attribute
¶
The latest a poll started after its slot, in milliseconds.
reconnects
class-attribute
instance-attribute
¶
Outages that ended with the analyzer reopened.
samples_dropped
class-attribute
instance-attribute
¶
Batches the overflow policy discarded, and batches still unread when the recording stopped.
samples_emitted
class-attribute
instance-attribute
¶
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
¶
Ticks skipped because a poll overran their slots.
target_total_samples
class-attribute
instance-attribute
¶
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.
ReconnectPolicy
dataclass
¶
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; |
__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
delay ¶
Recording
dataclass
¶
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. |
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: |
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
|
names
|
Sequence[str] | None
|
The analyzers to record; all of the source's when |
None
|
overflow
|
OverflowPolicy
|
What to do when |
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
|
FujiConnectionError
|
raised on leaving the block when a connection failure ended the recording. |
Source code in src/fujilib/streaming/recorder.py
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 ¶
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")