Skip to content

fujilib.sync

Blocking twins of the async API (design §7.3): Fuji.open, SyncAnalyzer and its SyncRemoteCalibration, discovery, recording and sinks, and the SyncPortal that runs the event loop in a background thread.

fujilib.sync.analyzer

The blocking analyzer facade: :class:SyncAnalyzer and :meth:Fuji.open (design §7.3).

Each method is a one-line blocking call of the :class:~fujilib.devices.analyzer.Analyzer method of the same name, with the same parameters and defaults (a parity test holds them together)::

from fujilib.sync import Fuji

with Fuji.open("COM8", channel_map={"CH3": "o2"}) as anz:
    print(anz.poll().channel("CH3"))

A remote manual calibration is a blocking context manager too (:class:SyncRemoteCalibration), entered and left on the portal's loop.

Fuji

The sync entry point: with Fuji.open(...) as anz:.

open staticmethod

open(
    port,
    *,
    profile=ZP_PROFILE,
    protocol=None,
    address=1,
    serial_settings=None,
    timeout=DEFAULTS.request_timeout_s,
    identify=True,
    channel_map=None,
    options=Capability.NONE,
    write_warn_per_minute=DEFAULTS.write_warn_per_minute,
    settle_after_reopen_s=DEFAULTS.settle_after_reopen_s,
    portal=None,
)

Open an analyzer for the with block; the arguments are :func:open_device's.

Without portal the analyzer gets a portal of its own, closed with it. A :class:~fujilib.transport.base.Transport passed as port must belong to portal's loop.

Source code in src/fujilib/sync/analyzer.py
@staticmethod
@contextmanager
def open(
    port: str | Transport,
    *,
    profile: DeviceProfile = ZP_PROFILE,
    protocol: ProtocolKind | str | None = None,
    address: int = 1,
    serial_settings: SerialSettings | None = None,
    timeout: float = DEFAULTS.request_timeout_s,
    identify: bool = True,
    channel_map: Mapping[ChannelId | str, Gas | str] | None = None,
    options: Capability = Capability.NONE,
    write_warn_per_minute: int = DEFAULTS.write_warn_per_minute,
    settle_after_reopen_s: float = DEFAULTS.settle_after_reopen_s,
    portal: SyncPortal | None = None,
) -> Generator[SyncAnalyzer]:
    """Open an analyzer for the ``with`` block; the arguments are :func:`open_device`'s.

    Without ``portal`` the analyzer gets a portal of its own, closed with
    it. A :class:`~fujilib.transport.base.Transport` passed as ``port``
    must belong to ``portal``'s loop.
    """
    with ExitStack() as stack:
        active = portal if portal is not None else stack.enter_context(SyncPortal())
        analyzer = active.call(
            open_device,
            port,
            profile=profile,
            protocol=protocol,
            address=address,
            serial_settings=serial_settings,
            timeout=timeout,
            identify=identify,
            channel_map=channel_map,
            options=options,
            write_warn_per_minute=write_warn_per_minute,
            settle_after_reopen_s=settle_after_reopen_s,
        )
        stack.callback(active.call, analyzer.close)
        yield SyncAnalyzer(analyzer, active)

SyncAnalyzer

SyncAnalyzer(analyzer, portal)

A blocking view of an :class:~fujilib.devices.analyzer.Analyzer, bound to a portal.

Wrap analyzer, whose loop is portal's; :meth:Fuji.open does this.

Source code in src/fujilib/sync/analyzer.py
def __init__(self, analyzer: Analyzer, portal: SyncPortal) -> None:
    """Wrap ``analyzer``, whose loop is ``portal``'s; :meth:`Fuji.open` does this."""
    self._anz = analyzer
    self._portal = portal

address property

address

:attr:Analyzer.address.

analyzer property

analyzer

The async analyzer; its coroutines must run on :attr:portal.

channels property

channels

:attr:Analyzer.channels.

info property

info

:attr:Analyzer.info.

last_frame property

last_frame

:attr:Analyzer.last_frame.

options property

options

:attr:Analyzer.options.

port property

port

:attr:Analyzer.port.

portal property

portal

The portal the analyzer's loop runs in.

protocol property

protocol

:attr:Analyzer.protocol.

session property

session

:attr:Analyzer.session.

apply_settings

apply_settings(
    document,
    *,
    confirm=False,
    any_analyzer=False,
    max_tier=SafetyTier.DANGEROUS,
    timeout=None,
)

Blocking :meth:Analyzer.apply_settings.

Source code in src/fujilib/sync/analyzer.py
def apply_settings(
    self,
    document: SettingsDocument | Mapping[str, object],
    *,
    confirm: bool = False,
    any_analyzer: bool = False,
    max_tier: SafetyTier = SafetyTier.DANGEROUS,
    timeout: float | None = None,
) -> ApplyReport:
    """Blocking :meth:`Analyzer.apply_settings`."""
    return self._portal.call(
        self._anz.apply_settings,
        document,
        confirm=confirm,
        any_analyzer=any_analyzer,
        max_tier=max_tier,
        timeout=timeout,
    )

calibration_status

calibration_status(*, timeout=None)

Blocking :meth:Analyzer.calibration_status.

Source code in src/fujilib/sync/analyzer.py
def calibration_status(self, *, timeout: float | None = None) -> CalibrationStatus:
    """Blocking :meth:`Analyzer.calibration_status`."""
    return self._portal.call(self._anz.calibration_status, timeout=timeout)

channel_status

channel_status(channel, *, timeout=None)

Blocking :meth:Analyzer.channel_status.

Source code in src/fujilib/sync/analyzer.py
def channel_status(
    self, channel: ChannelId | str, *, timeout: float | None = None
) -> ChannelStatus:
    """Blocking :meth:`Analyzer.channel_status`."""
    return self._portal.call(self._anz.channel_status, channel, timeout=timeout)

close

close()

Blocking :meth:Analyzer.close.

Source code in src/fujilib/sync/analyzer.py
def close(self) -> None:
    """Blocking :meth:`Analyzer.close`."""
    self._portal.call(self._anz.close)

diff_settings

diff_settings(
    document, *, any_analyzer=False, timeout=None
)

Blocking :meth:Analyzer.diff_settings.

Source code in src/fujilib/sync/analyzer.py
def diff_settings(
    self,
    document: SettingsDocument | Mapping[str, object],
    *,
    any_analyzer: bool = False,
    timeout: float | None = None,
) -> SettingsDiff:
    """Blocking :meth:`Analyzer.diff_settings`."""
    return self._portal.call(
        self._anz.diff_settings, document, any_analyzer=any_analyzer, timeout=timeout
    )

identify

identify(*, channel_map=None, timeout=None)

Blocking :meth:Analyzer.identify.

Source code in src/fujilib/sync/analyzer.py
def identify(
    self,
    *,
    channel_map: Mapping[ChannelId | str, Gas | str] | None = None,
    timeout: float | None = None,
) -> DeviceInfo:
    """Blocking :meth:`Analyzer.identify`."""
    return self._portal.call(self._anz.identify, channel_map=channel_map, timeout=timeout)

manual_calibration

manual_calibration(
    plan,
    *,
    gas,
    confirm=False,
    rule=None,
    adc=False,
    interval=0.5,
    key_timeout=2.0,
    run_timeout=30.0,
    cleanup_timeout=30.0,
)

Blocking :meth:Analyzer.manual_calibration; use it in a with block.

Source code in src/fujilib/sync/analyzer.py
def manual_calibration(
    self,
    plan: ManualCalibrationPlan,
    *,
    gas: CalibrationGas | Mapping[ChannelId | str, CalibrationGas],
    confirm: bool = False,
    rule: SteadinessRule | None = None,
    adc: bool = False,
    interval: float = 0.5,
    key_timeout: float = 2.0,
    run_timeout: float = 30.0,
    cleanup_timeout: float = 30.0,
) -> SyncRemoteCalibration:
    """Blocking :meth:`Analyzer.manual_calibration`; use it in a ``with`` block."""
    run = self._anz.manual_calibration(
        plan,
        gas=gas,
        confirm=confirm,
        rule=rule,
        adc=adc,
        interval=interval,
        key_timeout=key_timeout,
        run_timeout=run_timeout,
        cleanup_timeout=cleanup_timeout,
    )
    return SyncRemoteCalibration(run, self._portal)

plan_auto_calibration

plan_auto_calibration(*, timeout=None)

Blocking :meth:Analyzer.plan_auto_calibration.

Source code in src/fujilib/sync/analyzer.py
def plan_auto_calibration(self, *, timeout: float | None = None) -> CalibrationPlan:
    """Blocking :meth:`Analyzer.plan_auto_calibration`."""
    return self._portal.call(self._anz.plan_auto_calibration, timeout=timeout)

plan_auto_zero_calibration

plan_auto_zero_calibration(*, timeout=None)

Blocking :meth:Analyzer.plan_auto_zero_calibration.

Source code in src/fujilib/sync/analyzer.py
def plan_auto_zero_calibration(self, *, timeout: float | None = None) -> CalibrationPlan:
    """Blocking :meth:`Analyzer.plan_auto_zero_calibration`."""
    return self._portal.call(self._anz.plan_auto_zero_calibration, timeout=timeout)

plan_manual_calibration

plan_manual_calibration(channel, kind, *, timeout=None)

Blocking :meth:Analyzer.plan_manual_calibration.

Source code in src/fujilib/sync/analyzer.py
def plan_manual_calibration(
    self,
    channel: ChannelId | str,
    kind: ManualCalibrationKind | str,
    *,
    timeout: float | None = None,
) -> ManualCalibrationPlan:
    """Blocking :meth:`Analyzer.plan_manual_calibration`."""
    return self._portal.call(self._anz.plan_manual_calibration, channel, kind, timeout=timeout)

poll

poll(*, detail=True, timeout=None)

Blocking :meth:Analyzer.poll.

Source code in src/fujilib/sync/analyzer.py
def poll(self, *, detail: bool = True, timeout: float | None = None) -> Frame:
    """Blocking :meth:`Analyzer.poll`."""
    return self._portal.call(self._anz.poll, detail=detail, timeout=timeout)

read_adc

read_adc(*, timeout=None)

Blocking :meth:Analyzer.read_adc.

Source code in src/fujilib/sync/analyzer.py
def read_adc(self, *, timeout: float | None = None) -> AdcValues:
    """Blocking :meth:`Analyzer.read_adc`."""
    return self._portal.call(self._anz.read_adc, timeout=timeout)

read_calibration_log

read_calibration_log(channel=None, *, timeout=None)

Blocking :meth:Analyzer.read_calibration_log.

Source code in src/fujilib/sync/analyzer.py
def read_calibration_log(
    self, channel: ChannelId | str | None = None, *, timeout: float | None = None
) -> tuple[CalibrationLogEntry, ...]:
    """Blocking :meth:`Analyzer.read_calibration_log`."""
    return self._portal.call(self._anz.read_calibration_log, channel, timeout=timeout)

read_channel

read_channel(channel, *, timeout=None)

Blocking :meth:Analyzer.read_channel.

Source code in src/fujilib/sync/analyzer.py
def read_channel(self, channel: ChannelId | str, *, timeout: float | None = None) -> Reading:
    """Blocking :meth:`Analyzer.read_channel`."""
    return self._portal.call(self._anz.read_channel, channel, timeout=timeout)

read_clock

read_clock(*, timeout=None)

Blocking :meth:Analyzer.read_clock.

Source code in src/fujilib/sync/analyzer.py
def read_clock(self, *, timeout: float | None = None) -> ClockReading:
    """Blocking :meth:`Analyzer.read_clock`."""
    return self._portal.call(self._anz.read_clock, timeout=timeout)

read_error_log

read_error_log(*, timeout=None)

Blocking :meth:Analyzer.read_error_log.

Source code in src/fujilib/sync/analyzer.py
def read_error_log(self, *, timeout: float | None = None) -> tuple[ErrorLogEntry, ...]:
    """Blocking :meth:`Analyzer.read_error_log`."""
    return self._portal.call(self._anz.read_error_log, timeout=timeout)

read_metadata

read_metadata(*, timeout=None)

Blocking :meth:Analyzer.read_metadata.

Source code in src/fujilib/sync/analyzer.py
def read_metadata(self, *, timeout: float | None = None) -> AnalyzerMetadata:
    """Blocking :meth:`Analyzer.read_metadata`."""
    return self._portal.call(self._anz.read_metadata, timeout=timeout)

read_parameter

read_parameter(name, *, alarm_targets=None, timeout=None)

Blocking :meth:Analyzer.read_parameter.

Source code in src/fujilib/sync/analyzer.py
def read_parameter(
    self,
    name: str,
    *,
    alarm_targets: Mapping[int, ChannelId | str] | None = None,
    timeout: float | None = None,
) -> RegisterValue:
    """Blocking :meth:`Analyzer.read_parameter`."""
    return self._portal.call(
        self._anz.read_parameter, name, alarm_targets=alarm_targets, timeout=timeout
    )

read_parameters

read_parameters(names, *, alarm_targets=None, timeout=None)

Blocking :meth:Analyzer.read_parameters.

Source code in src/fujilib/sync/analyzer.py
def read_parameters(
    self,
    names: Iterable[str],
    *,
    alarm_targets: Mapping[int, ChannelId | str] | None = None,
    timeout: float | None = None,
) -> Mapping[str, RegisterValue]:
    """Blocking :meth:`Analyzer.read_parameters`."""
    return self._portal.call(
        self._anz.read_parameters, names, alarm_targets=alarm_targets, timeout=timeout
    )

read_ranges

read_ranges(*, timeout=None)

Blocking :meth:Analyzer.read_ranges.

Source code in src/fujilib/sync/analyzer.py
def read_ranges(self, *, timeout: float | None = None) -> tuple[RangeInfo, ...]:
    """Blocking :meth:`Analyzer.read_ranges`."""
    return self._portal.call(self._anz.read_ranges, timeout=timeout)

read_settings

read_settings(*, alarm_targets=None, timeout=None)

Blocking :meth:Analyzer.read_settings.

Source code in src/fujilib/sync/analyzer.py
def read_settings(
    self,
    *,
    alarm_targets: Mapping[int, ChannelId | str] | None = None,
    timeout: float | None = None,
) -> Mapping[str, RegisterValue]:
    """Blocking :meth:`Analyzer.read_settings`."""
    return self._portal.call(
        self._anz.read_settings, alarm_targets=alarm_targets, timeout=timeout
    )

reopen

reopen(*, timeout=None)

Blocking :meth:Analyzer.reopen.

Source code in src/fujilib/sync/analyzer.py
def reopen(self, *, timeout: float | None = None) -> DeviceInfo:
    """Blocking :meth:`Analyzer.reopen`."""
    return self._portal.call(self._anz.reopen, timeout=timeout)

reprobe

reprobe(capability, *, timeout=None)

Blocking :meth:Analyzer.reprobe.

Source code in src/fujilib/sync/analyzer.py
def reprobe(self, capability: Capability, *, timeout: float | None = None) -> Availability:
    """Blocking :meth:`Analyzer.reprobe`."""
    return self._portal.call(self._anz.reprobe, capability, timeout=timeout)

return_to_measurement

return_to_measurement(*, confirm=False, timeout=None)

Blocking :meth:Analyzer.return_to_measurement.

Source code in src/fujilib/sync/analyzer.py
def return_to_measurement(
    self, *, confirm: bool = False, timeout: float | None = None
) -> CommandResult:
    """Blocking :meth:`Analyzer.return_to_measurement`."""
    return self._portal.call(self._anz.return_to_measurement, confirm=confirm, timeout=timeout)

set_calibration_gas

set_calibration_gas(
    channel,
    range_number,
    kind,
    value,
    *,
    unit,
    confirm=False,
    timeout=None,
)

Blocking :meth:Analyzer.set_calibration_gas.

Source code in src/fujilib/sync/analyzer.py
def set_calibration_gas(
    self,
    channel: ChannelId | str,
    range_number: int,
    kind: str,
    value: float | str,
    *,
    unit: Unit | str,
    confirm: bool = False,
    timeout: float | None = None,
) -> WriteResult:
    """Blocking :meth:`Analyzer.set_calibration_gas`."""
    return self._portal.call(
        self._anz.set_calibration_gas,
        channel,
        range_number,
        kind,
        value,
        unit=unit,
        confirm=confirm,
        timeout=timeout,
    )

set_hold_mode

set_hold_mode(mode, *, confirm=False, timeout=None)

Blocking :meth:Analyzer.set_hold_mode.

Source code in src/fujilib/sync/analyzer.py
def set_hold_mode(
    self, mode: HoldMode | str, *, confirm: bool = False, timeout: float | None = None
) -> WriteResult:
    """Blocking :meth:`Analyzer.set_hold_mode`."""
    return self._portal.call(self._anz.set_hold_mode, mode, confirm=confirm, timeout=timeout)

set_hold_value

set_hold_value(
    channel, percent_fs, *, confirm=False, timeout=None
)

Blocking :meth:Analyzer.set_hold_value.

Source code in src/fujilib/sync/analyzer.py
def set_hold_value(
    self,
    channel: ChannelId | str,
    percent_fs: int,
    *,
    confirm: bool = False,
    timeout: float | None = None,
) -> WriteResult:
    """Blocking :meth:`Analyzer.set_hold_value`."""
    return self._portal.call(
        self._anz.set_hold_value, channel, percent_fs, confirm=confirm, timeout=timeout
    )

set_output_hold

set_output_hold(enabled, *, confirm=False, timeout=None)

Blocking :meth:Analyzer.set_output_hold.

Source code in src/fujilib/sync/analyzer.py
def set_output_hold(
    self, enabled: bool, *, confirm: bool = False, timeout: float | None = None
) -> WriteResult:
    """Blocking :meth:`Analyzer.set_output_hold`."""
    return self._portal.call(
        self._anz.set_output_hold, enabled, confirm=confirm, timeout=timeout
    )

set_range

set_range(
    channel, range_number, *, confirm=False, timeout=None
)

Blocking :meth:Analyzer.set_range.

Source code in src/fujilib/sync/analyzer.py
def set_range(
    self,
    channel: ChannelId | str,
    range_number: int,
    *,
    confirm: bool = False,
    timeout: float | None = None,
) -> WriteResult:
    """Blocking :meth:`Analyzer.set_range`."""
    return self._portal.call(
        self._anz.set_range, channel, range_number, confirm=confirm, timeout=timeout
    )

set_range_method

set_range_method(
    channel, method, *, confirm=False, timeout=None
)

Blocking :meth:Analyzer.set_range_method.

Source code in src/fujilib/sync/analyzer.py
def set_range_method(
    self,
    channel: ChannelId | str,
    method: RangeMethod | str,
    *,
    confirm: bool = False,
    timeout: float | None = None,
) -> WriteResult:
    """Blocking :meth:`Analyzer.set_range_method`."""
    return self._portal.call(
        self._anz.set_range_method, channel, method, confirm=confirm, timeout=timeout
    )

set_response_time

set_response_time(
    target, seconds, *, confirm=False, timeout=None
)

Blocking :meth:Analyzer.set_response_time.

Source code in src/fujilib/sync/analyzer.py
def set_response_time(
    self,
    target: ChannelId | str,
    seconds: int,
    *,
    confirm: bool = False,
    timeout: float | None = None,
) -> WriteResult:
    """Blocking :meth:`Analyzer.set_response_time`."""
    return self._portal.call(
        self._anz.set_response_time, target, seconds, confirm=confirm, timeout=timeout
    )

snapshot

snapshot(*, name=None)

Blocking :meth:Analyzer.snapshot.

Source code in src/fujilib/sync/analyzer.py
def snapshot(self, *, name: str | None = None) -> FujiDeviceSnapshot:
    """Blocking :meth:`Analyzer.snapshot`."""
    return self._portal.call(self._anz.snapshot, name=name)

start_auto_calibration

start_auto_calibration(*, confirm=False, timeout=None)

Blocking :meth:Analyzer.start_auto_calibration.

Source code in src/fujilib/sync/analyzer.py
def start_auto_calibration(
    self, *, confirm: bool = False, timeout: float | None = None
) -> CommandResult:
    """Blocking :meth:`Analyzer.start_auto_calibration`."""
    return self._portal.call(self._anz.start_auto_calibration, confirm=confirm, timeout=timeout)

start_auto_zero_calibration

start_auto_zero_calibration(*, confirm=False, timeout=None)

Blocking :meth:Analyzer.start_auto_zero_calibration.

Source code in src/fujilib/sync/analyzer.py
def start_auto_zero_calibration(
    self, *, confirm: bool = False, timeout: float | None = None
) -> CommandResult:
    """Blocking :meth:`Analyzer.start_auto_zero_calibration`."""
    return self._portal.call(
        self._anz.start_auto_zero_calibration, confirm=confirm, timeout=timeout
    )

start_blowback

start_blowback(*, confirm=False, timeout=None)

Blocking :meth:Analyzer.start_blowback.

Source code in src/fujilib/sync/analyzer.py
def start_blowback(
    self, *, confirm: bool = False, timeout: float | None = None
) -> CommandResult:
    """Blocking :meth:`Analyzer.start_blowback`."""
    return self._portal.call(self._anz.start_blowback, confirm=confirm, timeout=timeout)

status

status(*, timeout=None)

Blocking :meth:Analyzer.status.

Source code in src/fujilib/sync/analyzer.py
def status(self, *, timeout: float | None = None) -> AnalyzerStatus:
    """Blocking :meth:`Analyzer.status`."""
    return self._portal.call(self._anz.status, timeout=timeout)

wait_for_calibration

wait_for_calibration(*, timeout, interval=2.0, since=None)

Blocking :meth:Analyzer.wait_for_calibration.

Source code in src/fujilib/sync/analyzer.py
def wait_for_calibration(
    self,
    *,
    timeout: float,
    interval: float = 2.0,
    since: CalibrationStatus | None = None,
) -> CalibrationWait:
    """Blocking :meth:`Analyzer.wait_for_calibration`."""
    return self._portal.call(
        self._anz.wait_for_calibration, timeout=timeout, interval=interval, since=since
    )

wait_for_manual_calibration

wait_for_manual_calibration(
    *, timeout, interval=0.5, adc=False
)

Blocking :meth:Analyzer.wait_for_manual_calibration.

Source code in src/fujilib/sync/analyzer.py
def wait_for_manual_calibration(
    self, *, timeout: float, interval: float = 0.5, adc: bool = False
) -> ManualCalibrationEvent:
    """Blocking :meth:`Analyzer.wait_for_manual_calibration`."""
    return self._portal.call(
        self._anz.wait_for_manual_calibration, timeout=timeout, interval=interval, adc=adc
    )

write_parameter

write_parameter(
    name, value, *, unit=None, confirm=False, timeout=None
)

Blocking :meth:Analyzer.write_parameter.

Source code in src/fujilib/sync/analyzer.py
def write_parameter(
    self,
    name: str,
    value: object,
    *,
    unit: Unit | str | None = None,
    confirm: bool = False,
    timeout: float | None = None,
) -> WriteResult:
    """Blocking :meth:`Analyzer.write_parameter`."""
    return self._portal.call(
        self._anz.write_parameter, name, value, unit=unit, confirm=confirm, timeout=timeout
    )

SyncRemoteCalibration

SyncRemoteCalibration(run, portal)

A blocking view of a :class:~fujilib.devices.keys.RemoteCalibration.

Its with block enters and leaves the run on the portal's loop, so the cleanup runs there however the block is left. Every call is interruptible (:meth:SyncPortal.call_interruptible): Ctrl-C cancels it on the loop and waits for it to finish before the block is left, so a key it was about to send is not sent afterwards, and Ctrl-C while entering still cleans up. progress callbacks are called on the loop's thread.

Wrap run, whose analyzer's loop is portal's.

Source code in src/fujilib/sync/analyzer.py
def __init__(self, run: RemoteCalibration, portal: SyncPortal) -> None:
    """Wrap ``run``, whose analyzer's loop is ``portal``'s."""
    self._run = run
    self._portal = portal

event property

event

:attr:RemoteCalibration.event.

gases property

gases

:attr:RemoteCalibration.gases.

keys property

keys

:attr:RemoteCalibration.keys.

observation property

observation

:attr:RemoteCalibration.observation.

plan property

plan

:attr:RemoteCalibration.plan.

result property

result

:attr:RemoteCalibration.result.

rule property

rule

:attr:RemoteCalibration.rule.

run property

run

The async run; its coroutines must run on the portal.

state property

state

:attr:RemoteCalibration.state.

steadiness property

steadiness

:attr:RemoteCalibration.steadiness.

calibrate

calibrate(*, confirm=False)

Blocking :meth:RemoteCalibration.calibrate.

Source code in src/fujilib/sync/analyzer.py
def calibrate(self, *, confirm: bool = False) -> ManualCalibrationEvent:
    """Blocking :meth:`RemoteCalibration.calibrate`."""
    return self._portal.call_interruptible(self._run.calibrate, confirm=confirm)

cancel

cancel()

Blocking :meth:RemoteCalibration.cancel.

Source code in src/fujilib/sync/analyzer.py
def cancel(self) -> ManualCalibrationEvent | None:
    """Blocking :meth:`RemoteCalibration.cancel`."""
    return self._portal.call_interruptible(self._run.cancel)

read

read(*, timeout=None)

Blocking :meth:RemoteCalibration.read.

Source code in src/fujilib/sync/analyzer.py
def read(self, *, timeout: float | None = None) -> SteadinessVerdict:
    """Blocking :meth:`RemoteCalibration.read`."""
    return self._portal.call_interruptible(self._run.read, timeout=timeout)

wait_steady

wait_steady(*, timeout=None, progress=None)

Blocking :meth:RemoteCalibration.wait_steady.

Source code in src/fujilib/sync/analyzer.py
def wait_steady(
    self,
    *,
    timeout: float | None = None,
    progress: Callable[[SteadinessVerdict], object] | None = None,
) -> SteadinessVerdict:
    """Blocking :meth:`RemoteCalibration.wait_steady`."""
    return self._portal.call_interruptible(
        self._run.wait_steady, timeout=timeout, progress=progress
    )

fujilib.sync.discovery

Blocking discovery (design §7.5).

find_devices

find_devices(
    *,
    ports=None,
    addresses=(1,),
    profiles=DEVICE_PROFILES,
    per_probe_timeout_s=0.3,
    identify=True,
    max_concurrency=8,
    portal=None,
)

Blocking :func:fujilib.devices.discovery.find_devices.

Runs on portal, or on a portal of its own when portal is None.

Source code in src/fujilib/sync/discovery.py
def find_devices(
    *,
    ports: Sequence[str] | None = None,
    addresses: Sequence[int] = (1,),
    profiles: Sequence[DeviceProfile] = DEVICE_PROFILES,
    per_probe_timeout_s: float = 0.3,
    identify: bool = True,
    max_concurrency: int = 8,
    portal: SyncPortal | None = None,
) -> list[DiscoveryResult]:
    """Blocking :func:`fujilib.devices.discovery.find_devices`.

    Runs on ``portal``, or on a portal of its own when ``portal`` is ``None``.
    """
    with SyncPortal() if portal is None else nullcontext(portal) as active:
        return active.call(
            _find_devices,
            ports=ports,
            addresses=addresses,
            profiles=profiles,
            per_probe_timeout_s=per_probe_timeout_s,
            identify=identify,
            max_concurrency=max_concurrency,
        )

fujilib.sync.recording

Blocking recording (design §7.3, §7.6).

:func:record runs the async recorder on a portal's loop and yields a :class:SyncRecording, iterated with a plain for loop::

from fujilib.sync import Fuji, PollSourceAdapter, SyncCsvSink, pipe, record

with (
    Fuji.open("COM8", channel_map={"CH3": "o2"}) as anz,
    record(PollSourceAdapter("zpa", anz), rate_hz=1.0, duration=60) as rec,
    SyncCsvSink("run.csv", portal=anz.portal) as sink,
):
    pipe(rec, sink)

Everything behaves as in :mod:fujilib.streaming.recorder; leaving the with block stops the recording, and raises the error that ended it, if one did.

PollSourceAdapter

PollSourceAdapter(name, device)

A :class:~fujilib.sync.analyzer.SyncAnalyzer as a poll source (unified API §E).

The blocking twin of :class:fujilib.streaming.poll_source.PollSourceAdapter; it brings the analyzer's portal to :func:record.

Publish device's polls under name.

Source code in src/fujilib/sync/recording.py
def __init__(self, name: str, device: SyncAnalyzer) -> None:
    """Publish ``device``'s polls under ``name``."""
    self._device = device
    self._source = AsyncPollSourceAdapter(name, device.analyzer)

device property

device

The wrapped analyzer.

name property

name

The name the analyzer's data is published under.

portal property

portal

The portal the analyzer's loop runs in.

source property

source

The async poll source, whose coroutines run on :attr:portal.

layout

layout(names=None)

:meth:fujilib.streaming.poll_source.PollSourceAdapter.layout.

Source code in src/fujilib/sync/recording.py
def layout(self, names: Sequence[str] | None = None) -> Mapping[str, SourceLayout]:
    """:meth:`fujilib.streaming.poll_source.PollSourceAdapter.layout`."""
    return self._source.layout(names)

poll

poll(names=None)

Blocking :meth:fujilib.streaming.poll_source.PollSourceAdapter.poll.

Source code in src/fujilib/sync/recording.py
def poll(self, names: Sequence[str] | None = None) -> Mapping[str, DeviceResult[Frame]]:
    """Blocking :meth:`fujilib.streaming.poll_source.PollSourceAdapter.poll`."""
    return self.portal.call(self._source.poll, names)

reconnect

reconnect(name)

Blocking :meth:fujilib.streaming.poll_source.PollSourceAdapter.reconnect.

Source code in src/fujilib/sync/recording.py
def reconnect(self, name: str) -> None:
    """Blocking :meth:`fujilib.streaming.poll_source.PollSourceAdapter.reconnect`."""
    self.portal.call(self._source.reconnect, name)

SyncRecording

SyncRecording(recording, portal)

A running recording, iterated with a plain for loop (unified API §I).

Wrap recording, which runs on portal's loop; :func:record does this.

Source code in src/fujilib/sync/recording.py
def __init__(self, recording: Recording[T], portal: SyncPortal) -> None:
    """Wrap ``recording``, which runs on ``portal``'s loop; :func:`record` does this."""
    self._recording = recording
    self._portal = portal

portal property

portal

The portal the recording runs in.

rate_hz property

rate_hz

The requested rate.

recording property

recording

The async recording, on :attr:portal.

stream property

stream

The batches, one per tick: this recording itself.

summary property

summary

The live counters.

pipe

pipe(source, sink, *, batch_size=64, flush_interval=1.0)

Blocking :func:fujilib.sinks.base.pipe, on the recording's portal.

sink is a blocking sink or an async one, open either way. Ctrl-C stops it: the pipe writes what it holds, then KeyboardInterrupt is raised here.

Source code in src/fujilib/sync/recording.py
def pipe(
    source: SyncRecording[Batch],
    sink: SyncSinkAdapter | SampleSink,
    *,
    batch_size: int = 64,
    flush_interval: float = 1.0,
) -> AcquisitionSummary:
    """Blocking :func:`fujilib.sinks.base.pipe`, on the recording's portal.

    ``sink`` is a blocking sink or an async one, open either way. Ctrl-C
    stops it: the pipe writes what it holds, then ``KeyboardInterrupt`` is
    raised here.
    """
    target = sink.async_sink if isinstance(sink, SyncSinkAdapter) else sink
    portal = source.portal
    scopes: list[anyio.CancelScope] = []

    async def run() -> AcquisitionSummary | None:
        with anyio.CancelScope() as scope:
            scopes.append(scope)
            return await async_pipe(
                source.recording, target, batch_size=batch_size, flush_interval=flush_interval
            )
        return None  # cancelled by Ctrl-C

    future = portal.start_task_soon(run)
    try:
        summary = _wait(future)
    except KeyboardInterrupt:

        def stop() -> None:
            for scope in scopes:
                scope.cancel()

        portal.run_in_loop(stop)
        with contextlib.suppress(Exception):
            _ = _wait(future)  # the pipe's last write
        raise
    assert summary is not None  # noqa: S101 - only a cancelled run returns None
    return summary

record

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

Blocking :func:fujilib.streaming.recorder.record.

source is a blocking :class:PollSourceAdapter, whose portal is used unless portal is given, or an async poll source with the portal its analyzers run on.

Raises:

Type Description
FujiValidationError

an async source without a portal, or an argument :func:~fujilib.streaming.recorder.record refuses.

FujiConnectionError

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

Source code in src/fujilib/sync/recording.py
@contextmanager
def record(
    source: PollSourceAdapter | 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,
    portal: SyncPortal | None = None,
) -> Generator[SyncRecording[Batch]]:
    """Blocking :func:`fujilib.streaming.recorder.record`.

    ``source`` is a blocking :class:`PollSourceAdapter`, whose portal is used
    unless ``portal`` is given, or an async poll source with the ``portal``
    its analyzers run on.

    Raises:
        FujiValidationError: an async source without a ``portal``, or an
            argument :func:`~fujilib.streaming.recorder.record` refuses.
        FujiConnectionError: on leaving the block, when a connection failure
            ended the recording.
    """
    if isinstance(source, PollSourceAdapter):
        async_source: PollSource = source.source
        active = portal if portal is not None else source.portal
    else:
        if portal is None:
            msg = "an async poll source needs the portal its analyzers run on"
            raise FujiValidationError(msg)
        async_source, active = source, portal
    manager = async_record(
        async_source,
        rate_hz=rate_hz,
        duration=duration,
        names=names,
        overflow=overflow,
        buffer_size=buffer_size,
        reconnect=reconnect,
    )
    with active.wrap_async_context_manager(manager) as recording:
        yield SyncRecording(recording, active)

fujilib.sync.sinks

Blocking sinks (design §7.3, §7.6).

Each is the async sink of the same name behind a :class:SyncSinkAdapter, which opens, writes and closes it on a portal's loop: the portal given, or one of its own for the with block. Pass the analyzer's portal (anz.portal) to share one loop.

SyncCsvSink

SyncCsvSink(path, *, channels=None, portal=None)

Bases: SyncSinkAdapter

Blocking :class:~fujilib.sinks.csv.CsvSink.

A sink for path; its columns are locked now if channels are given.

Source code in src/fujilib/sync/sinks.py
def __init__(
    self,
    path: str | PathLike[str],
    *,
    channels: Iterable[ChannelId] | None = None,
    portal: SyncPortal | None = None,
) -> None:
    """A sink for ``path``; its columns are locked now if ``channels`` are given."""
    super().__init__(CsvSink(path, channels=channels), portal=portal)

SyncInMemorySink

SyncInMemorySink(*, channels=None, portal=None)

Bases: SyncSinkAdapter

Blocking :class:~fujilib.sinks.memory.InMemorySink.

An empty sink; its columns are locked now if channels are given.

Source code in src/fujilib/sync/sinks.py
def __init__(
    self, *, channels: Iterable[ChannelId] | None = None, portal: SyncPortal | None = None
) -> None:
    """An empty sink; its columns are locked now if ``channels`` are given."""
    self._memory = InMemorySink(channels=channels)
    super().__init__(self._memory, portal=portal)

samples property

samples

The samples written, in order.

rows

rows()

The rows written, under the sink's locked columns.

Source code in src/fujilib/sync/sinks.py
def rows(self) -> list[dict[str, Scalar]]:
    """The rows written, under the sink's locked columns."""
    return self._memory.rows()

SyncParquetSink

SyncParquetSink(
    path,
    *,
    channels=None,
    compression="zstd",
    row_group_size=DEFAULT_ROW_GROUP_SIZE,
    metadata=None,
    portal=None,
)

Bases: SyncSinkAdapter

Blocking :class:~fujilib.sinks.parquet.ParquetSink.

A sink for path; the arguments are :class:~fujilib.sinks.parquet.ParquetSink's.

Source code in src/fujilib/sync/sinks.py
def __init__(
    self,
    path: str | PathLike[str],
    *,
    channels: Iterable[ChannelId] | None = None,
    compression: Compression = "zstd",
    row_group_size: int = DEFAULT_ROW_GROUP_SIZE,
    metadata: Mapping[str, str] | None = None,
    portal: SyncPortal | None = None,
) -> None:
    """A sink for ``path``; the arguments are :class:`~fujilib.sinks.parquet.ParquetSink`'s."""
    sink = ParquetSink(
        path,
        channels=channels,
        compression=compression,
        row_group_size=row_group_size,
        metadata=metadata,
    )
    super().__init__(sink, portal=portal)

SyncSinkAdapter

SyncSinkAdapter(sink, *, portal=None)

A blocking view of an async :class:~fujilib.sinks.base.SampleSink.

Wrap sink; without portal it gets a portal of its own when opened.

Source code in src/fujilib/sync/sinks.py
def __init__(self, sink: SampleSink, *, portal: SyncPortal | None = None) -> None:
    """Wrap ``sink``; without ``portal`` it gets a portal of its own when opened."""
    self._sink = sink
    self._portal = portal
    self._owns_portal = portal is None
    self._stack = ExitStack()
    self._closed = False

async_sink property

async_sink

The async sink.

close

close()

Blocking close; a portal of the sink's own stops with it. Again is a no-op.

Raises:

Type Description
FujiSinkError

the shared portal stopped before the sink was closed, so the sink could not finish (a Parquet file would be unreadable).

Source code in src/fujilib/sync/sinks.py
def close(self) -> None:
    """Blocking ``close``; a portal of the sink's own stops with it. Again is a no-op.

    Raises:
        FujiSinkError: the shared portal stopped before the sink was closed, so
            the sink could not finish (a Parquet file would be unreadable).
    """
    if self._closed:
        return
    self._closed = True
    try:
        portal = self._portal
        if portal is None:
            return  # never opened
        if not portal.running:
            msg = "the sink's portal stopped before the sink was closed; close it first"
            raise FujiSinkError(msg)
        portal.call(self._sink.close)
    finally:
        self._stack.close()

open

open()

Blocking open; if it fails, a portal of the sink's own stops again.

Source code in src/fujilib/sync/sinks.py
def open(self) -> None:
    """Blocking ``open``; if it fails, a portal of the sink's own stops again."""
    try:
        self._active().call(self._sink.open)
    except BaseException:
        self._stack.close()
        if self._owns_portal:
            self._portal = None
        raise

write_many

write_many(samples)

Blocking write_many.

Source code in src/fujilib/sync/sinks.py
def write_many(self, samples: Sequence[Sample]) -> None:
    """Blocking ``write_many``."""
    self._active().call(self._sink.write_many, samples)

fujilib.sync.portal

The blocking portal behind the sync facade (design §7.3).

:class:SyncPortal runs an AnyIO event loop in a background thread (:func:anyio.from_thread.start_blocking_portal) and calls coroutines on it from ordinary code. A portal is used once: its with block starts the loop and ends it.

A single exception raised inside a task group reaches AnyIO wrapped in an exception group; :meth:SyncPortal.call unwraps a group of one, so callers catch the :class:~fujilib.errors.FujiError subclass itself.

SyncPortal

SyncPortal(*, backend='asyncio')

An event loop in a background thread, for blocking calls into the async core.

Example::

with SyncPortal() as portal:
    analyzer = portal.call(open_device, "COM8")

Prepare a portal on AnyIO backend ("asyncio" or "trio").

Source code in src/fujilib/sync/portal.py
def __init__(self, *, backend: str = "asyncio") -> None:
    """Prepare a portal on AnyIO ``backend`` (``"asyncio"`` or ``"trio"``)."""
    self._backend_name = backend
    self._cm: AbstractContextManager[BlockingPortal] | None = None
    self._portal: BlockingPortal | None = None
    self._entered = False

running property

running

Whether the portal's loop is running.

__enter__

__enter__()

Start the loop.

Raises:

Type Description
RuntimeError

the portal was used before.

Source code in src/fujilib/sync/portal.py
def __enter__(self) -> Self:
    """Start the loop.

    Raises:
        RuntimeError: the portal was used before.
    """
    if self._entered:
        msg = "a SyncPortal cannot be used again after it has been closed"
        raise RuntimeError(msg)
    self._entered = True
    cm = start_blocking_portal(self._backend_name)
    self._portal = cm.__enter__()
    self._cm = cm
    return self

__exit__

__exit__(exc_type, exc, tb)

Stop the loop and join its thread.

Source code in src/fujilib/sync/portal.py
def __exit__(
    self,
    exc_type: type[BaseException] | None,
    exc: BaseException | None,
    tb: TracebackType | None,
) -> None:
    """Stop the loop and join its thread."""
    cm, self._cm, self._portal = self._cm, None, None
    if cm is not None:
        cm.__exit__(exc_type, exc, tb)

call

call(func, *args, **kwargs)

Run func(*args, **kwargs) on the portal's loop and wait for the result.

Raises:

Type Description
RuntimeError

the portal is not running.

Source code in src/fujilib/sync/portal.py
def call[**P, T](self, func: Callable[P, Awaitable[T]], *args: P.args, **kwargs: P.kwargs) -> T:
    """Run ``func(*args, **kwargs)`` on the portal's loop and wait for the result.

    Raises:
        RuntimeError: the portal is not running.
    """
    portal = self._running()
    try:
        return portal.call(partial(func, *args, **kwargs))
    except BaseExceptionGroup as group:
        unwrapped = unwrap(group)
        if unwrapped is group:
            raise
    # Raised outside the handler: the group is hidden, and the error keeps
    # its own cause and context (errors.py).
    raise unwrapped

call_interruptible

call_interruptible(func, *args, **kwargs)

As :meth:call, but Ctrl-C cancels func on the loop before it is raised here.

:meth:call leaves the coroutine running on the loop when Ctrl-C interrupts the wait for it. Here Ctrl-C cancels it there and waits for it to finish, whatever it does on its way out, before KeyboardInterrupt is raised. So nothing it was about to do happens afterwards.

Raises:

Type Description
RuntimeError

the portal is not running.

Source code in src/fujilib/sync/portal.py
def call_interruptible[**P, T](
    self, func: Callable[P, Awaitable[T]], *args: P.args, **kwargs: P.kwargs
) -> T:
    """As :meth:`call`, but Ctrl-C cancels ``func`` on the loop before it is raised here.

    :meth:`call` leaves the coroutine running on the loop when Ctrl-C
    interrupts the wait for it. Here Ctrl-C cancels it there and waits for
    it to finish, whatever it does on its way out, before
    ``KeyboardInterrupt`` is raised. So nothing it was about to do happens
    afterwards.

    Raises:
        RuntimeError: the portal is not running.
    """
    portal = self._running()
    scopes: list[anyio.CancelScope] = []

    async def run() -> T | None:
        with anyio.CancelScope() as scope:
            scopes.append(scope)
            return await func(*args, **kwargs)
        return None  # cancelled by Ctrl-C

    future = portal.start_task_soon(run)
    try:
        result = _wait(future)
    except KeyboardInterrupt:

        def stop() -> None:
            for scope in scopes:
                scope.cancel()

        portal.call(stop)
        with contextlib.suppress(Exception):
            _wait(future)
        raise
    except BaseExceptionGroup as group:
        unwrapped = unwrap(group)
        if unwrapped is group:
            raise
    else:
        return result  # type: ignore[return-value]  # only a cancelled run returns None
    raise unwrapped

run_in_loop

run_in_loop(func)

Run the plain function func in the loop's thread and return its result.

Raises:

Type Description
RuntimeError

the portal is not running.

Source code in src/fujilib/sync/portal.py
def run_in_loop[T](self, func: Callable[[], T]) -> T:
    """Run the plain function ``func`` in the loop's thread and return its result.

    Raises:
        RuntimeError: the portal is not running.
    """
    return self._running().call(func)

start_task_soon

start_task_soon(func)

Start func() on the portal's loop; its future cancels the task when cancelled.

Raises:

Type Description
RuntimeError

the portal is not running.

Source code in src/fujilib/sync/portal.py
def start_task_soon[T](self, func: Callable[[], Awaitable[T]]) -> Future[T]:
    """Start ``func()`` on the portal's loop; its future cancels the task when cancelled.

    Raises:
        RuntimeError: the portal is not running.
    """
    return self._running().start_task_soon(func)

wrap_async_context_manager

wrap_async_context_manager(acm)

Present an async context manager, entered and exited on the portal's loop, as a sync one.

Raises:

Type Description
RuntimeError

the portal is not running.

Source code in src/fujilib/sync/portal.py
def wrap_async_context_manager[T](
    self, acm: AbstractAsyncContextManager[T]
) -> AbstractContextManager[T]:
    """Present an async context manager, entered and exited on the portal's loop, as a sync one.

    Raises:
        RuntimeError: the portal is not running.
    """
    return self._running().wrap_async_context_manager(acm)