diff --git a/.bumpversion.toml b/.bumpversion.toml index 36d463a..37ede3b 100644 --- a/.bumpversion.toml +++ b/.bumpversion.toml @@ -1,5 +1,5 @@ [tool.bumpversion] -current_version = "2.6.1" +current_version = "2.7.0" commit = false tag = false diff --git a/.gitignore b/.gitignore index 7155f3a..1af0b6c 100644 --- a/.gitignore +++ b/.gitignore @@ -1,6 +1,11 @@ # IDE settings .idea + +docs/*.pdf +**/*handoff.md + + # ci artefacts pytest*.xml diff --git a/examples/owh9830.rst b/examples/owh9830.rst new file mode 100644 index 0000000..04c2f78 --- /dev/null +++ b/examples/owh9830.rst @@ -0,0 +1,173 @@ +OWON OWH9830 power meter +------------------------ + +Connect directly to the meter's RS-232 port with SCPI selected and matching +baud rate (default 115200, 8N1):: + + from scpi.devices.owh9830 import serial + from scpi.wrapper import AIOWrapper + + dev = AIOWrapper(serial("/dev/tty.usbserial-A92TRYDJ")) + try: + volts = dev.measure_voltage("1A") + amps = dev.measure_current("1A") + degrees = dev.measure_phase_angle("1A") + watts = dev.measure_real_power("1A") + var = dev.measure_reactive_power("1A") + harmonics = dev.measure_harmonics("1A", max_order=7) + finally: + dev.quit() + +Measurements return ``Decimal`` values. Elements are ``1A``, ``1B``, ``1C`` +and ``1sigma`` (case insensitive). Harmonics support individual phases only; +``voltage`` and ``current`` tuples contain RMS amplitudes in volts and amps, +starting with the fundamental at index zero, through ``max_order`` (1--63). +Harmonic reads enable harmonic measurement mode and leave it enabled so data +keeps updating; scalar measurements also work in this mode. The first read waits +one update period after switching modes or requesting more than ten orders. +Orders 1--10 use one range command and two list queries. Reads update the display +order range. Higher orders use two queries per order because firmware V1.2.0 +truncates long list replies and can omit separators. Unavailable individual +readings (``----``) are returned as ``None`` +in their order's position. Reads do not guarantee an atomic snapshot. + +``OWH9830(transport)`` also accepts existing transports/protocols for async use. +Commands are spaced by more than 100 ms. Generic ``SYST:ERR?`` checks are +disabled because this command is unsupported. Blank or malformed measurements +raise ``ValueError``; if replies are blank, check the front-panel LOCAL state. +Use one device instance per meter. The interactive example is:: + + uv run --locked python examples/owh9830_serial.py /dev/tty.usbserial-A92TRYDJ --max-order 7 + +Held measurement sets +--------------------- + +Use ``measure_snapshot`` for voltage, current, phase angle and real/reactive power +from the same held set of meter-reported values:: + + single = dev.measure_snapshot("1A") + all_phases = dev.measure_snapshot("1A-C") + print(single["1A"]["voltage"]) + print(all_phases["1B"]["phase_angle"]) + +Use ``measure_harmonic_snapshot`` to read scalars and harmonics up to 10th order +(1--10) under the same held set:: + + harmonic_snap = dev.measure_harmonic_snapshot("1A", max_order=10) + print(harmonic_snap["1A"]["harmonics"]["voltage"]) + +The result always maps element names to dictionaries containing ``voltage`` (V), +``current`` (A), ``phase_angle`` (degrees), ``real_power`` (W) and +``reactive_power`` (var), and optionally ``harmonics`` (tuples of ``voltage`` and +``current`` harmonics). Single ``1A``, ``1B``, ``1C`` and ``1sigma`` selections +are supported for scalar snapshots, and ``1A``, ``1B``, ``1C`` and ``1A-C`` for +harmonic snapshots. Inputs are case insensitive; returned names are canonical. +Values are ``Decimal`` or ``None`` for unavailable readings (``----``), +including phase/power on unconnected phases. + +The method configures the numeric display page once per selection (five slots +in 8ITEM or fifteen in 16ITEM) and harmonic order range, then reuses it for +repeated polls. It queries the current HOLD state, enables HOLD if necessary, +reads the configured slots and harmonic lists, and restores HOLD in a ``finally`` +block. It waits one configured update period plus 110 ms after setup or its +previous HOLD release so rapid polling permits new frames to update. An existing +HOLD ON stays ON and returns its already-held values for the same selection. +Changing selection while HOLD is already ON raises ``ValueError``; release HOLD +first. The entire sequence excludes other I/O through the same device instance. +Measurement mode is preserved or switched to ``HARMONIC`` when harmonics are +requested. + +The numeric display configuration remains selected after the call. This API owns +those slots while polling; changing them on the front panel or through another +protocol/transport instance requires a new device instance before polling again. +Numeric, harmonic, and update-rate commands sent through ``dev.command`` +invalidate the cached selection. +The method does not change update rate or guarantee new data on every poll. +HOLD stabilizes reported readings; simultaneous ADC acquisition across channels +has not been established. + +The interactive all-phase and harmonic examples are:: + + uv run --locked python examples/owh9830_serial.py /dev/tty.usbserial-A92TRYDJ --element 1A-C + uv run --locked python examples/owh9830_serial.py /dev/tty.usbserial-A92TRYDJ --harmonics + +Streaming snapshots +------------------- + +Use ``measure_snapshots`` (or ``stream_snapshots``) to stream snapshot measurements +as fast as possible via an async iterator without toggling HOLD on each frame:: + + async for snapshot in dev.measure_snapshots("1A"): + print(snapshot["1A"]["voltage"], snapshot["1A"]["current"]) + +Use ``measure_harmonic_snapshots`` (or ``stream_harmonic_snapshots``) to stream +snapshots including harmonics up to order 10:: + + async for snapshot in dev.measure_harmonic_snapshots("1A", max_order=10): + print(snapshot["1A"]["voltage"], snapshot["1A"]["harmonics"]["voltage"]) + +Synchronous code using ``AIOWrapper`` supports streaming using a standard ``for`` loop:: + + for snapshot in dev.measure_snapshots("1A"): + print(snapshot["1A"]) + + for snapshot in dev.measure_harmonic_snapshots("1A"): + print(snapshot["1A"]["harmonics"]) + +The iterator configures the numeric display page slots and harmonic settings +once at stream start. While the iterator is active, commands that modify +instrument setup (such as harmonic measurements, single snapshots, or ``:NUM``, +``:RATE``, ``*RST``, ``:DISP:MOD``, ``:HOLD``, and ``:HARM:ORD`` configuration +commands) refuse to execute and raise ``RuntimeError`` to prevent configuration +conflicts. Read-only queries (e.g. ``measure_voltage`` or ``ask``) execute +safely between polls. + +The interactive examples are:: + + uv run --locked python examples/owh9830_serial.py /dev/tty.usbserial-A92TRYDJ --stream 5 + uv run --locked python examples/owh9830_serial.py /dev/tty.usbserial-A92TRYDJ --stream 5 --harmonics + uv run --locked python examples/owh9830_async_stream.py /dev/tty.usbserial-A92TRYDJ --count 5 --harmonics 10 + +Batching experiments +-------------------- + +Tested on 2026-10-02 with OWH9830 firmware V1.2.0, RS-232 115200 8N1, +0.5-second update period, a heater on 1A and unconnected 1B/1C. All experiment +display configurations and HOLD states were restored afterward. + +Compound queries separated by semicolons returned only the first voltage +response, even when current, phase and both powers followed. Numeric bulk reads +returned correctly ordered values, with no meaningful extra latency from powers: + +======================== ================ ================== +Selection Quantities Median reply time +======================== ================ ================== +1A V/I/phase 23.50 ms +1A V/I/phase/P/Q 23.68 ms +1A-C V/I/phase 33.48 ms +1A-C V/I/phase/P/Q 34.10 ms +======================== ================ ================== + +These are five-sample command-to-complete-response medians, excluding setup, +command-spacing waits and HOLD operations. A separate ten-sample comparison of +15-value reads measured 34.63 ms in NORM mode and 34.10 ms in HARMONIC mode; +normal mode provided no clear benefit. Numeric data continued updating in both. +HOLD ON produced identical repeated bulk readings and matching direct readings; +HOLD OFF resumed updates. Bulk phase readings have more digits than the direct +phase query, which reports one decimal place. + +Repeated ``measure_snapshot`` calls, including HOLD and pacing, measured about +833 ms per call at the 0.5-second update period for both 1A and 1A-C, excluding +initial configuration. Five consecutive calls returned five distinct sets for +each selection. Temporarily selecting a 0.1-second update period reduced the +all-phase median to 444 ms, but five calls returned only two distinct sets. +Faster polling therefore does not establish a higher acquisition rate. + +Stress tests confirmed a 128-byte output limit, including LF, for numeric bulk +replies as well as harmonic lists. Truncation can leave a plausible partial number, +so field counts alone cannot detect it. The snapshot method treats a response at +that limit as suspect and reads each selected measurement using direct +``:MEAS:...`` queries under the same HOLD. Indexed numeric-slot queries disagreed +with the bulk data in a stress test and are therefore avoided. The fallback adds +round trips and uses the direct phase query's one-decimal precision while +preserving the held set. diff --git a/examples/owh9830_async_stream.py b/examples/owh9830_async_stream.py new file mode 100644 index 0000000..9b3fc8e --- /dev/null +++ b/examples/owh9830_async_stream.py @@ -0,0 +1,60 @@ +#!/usr/bin/env python3 +"""Async streaming of OWH9830 snapshot measurements without AIOWrapper.""" + +import argparse +import asyncio +import sys + +from scpi.devices.owh9830 import serial + + +async def main() -> None: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("port", help="serial port path or URL (e.g. /dev/tty.usbserial-A92TRYDJ)") + parser.add_argument("--baudrate", type=int, default=115200, help="serial baudrate (default: 115200)") + parser.add_argument( + "--element", choices=("1A", "1B", "1C", "1sigma", "1A-C"), default="1A", help="element to stream" + ) + parser.add_argument( + "--harmonics", + type=int, + choices=range(0, 11), + default=0, + help="include harmonics up to order (1..10, default 0 for disabled)", + ) + parser.add_argument( + "--count", type=int, default=0, help="number of snapshots to read (0 or negative for continuous)" + ) + args = parser.parse_args() + + dev = serial(args.port, baudrate=args.baudrate) + try: + identity = await dev.identify() + print(f"Connected to: {identity}") + print( + f"Streaming snapshots for {args.element}" + f"{f' with harmonics 1..{args.harmonics}' if args.harmonics else ''} (Ctrl-C to stop)..." + ) + + stream = ( + dev.measure_harmonic_snapshots(args.element, max_order=args.harmonics) + if args.harmonics > 0 + else dev.measure_snapshots(args.element) + ) + received = 0 + async for snapshot in stream: + received += 1 + print(f"[{received}]", snapshot) + if args.count > 0 and received >= args.count: + break + except KeyboardInterrupt: + print("\nStreaming stopped by user.") + finally: + await dev.quit() + + +if __name__ == "__main__": + try: + asyncio.run(main()) + except KeyboardInterrupt: + sys.exit(0) diff --git a/examples/owh9830_serial.py b/examples/owh9830_serial.py new file mode 100644 index 0000000..a6f11c6 --- /dev/null +++ b/examples/owh9830_serial.py @@ -0,0 +1,49 @@ +#!/usr/bin/env python3 +"""Interactive OWH9830 over direct RS-232, e.g. /dev/tty.usbserial-A92TRYDJ.""" + +import argparse +import atexit +import os + +from scpi.devices.owh9830 import serial +from scpi.wrapper import AIOWrapper + +if __name__ == "__main__": + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("port") + parser.add_argument("--baudrate", type=int, default=115200) + parser.add_argument("--element", choices=("1A", "1B", "1C", "1sigma", "1A-C"), default="1A") + parser.add_argument("--max-order", type=int, choices=range(1, 64), default=7) + parser.add_argument("--harmonics", action="store_true", help="include harmonics (1..10) in snapshot measurements") + parser.add_argument( + "--stream", type=int, nargs="?", const=-1, default=0, help="stream snapshots (count or continuous)" + ) + args = parser.parse_args() + dev = AIOWrapper(serial(args.port, baudrate=args.baudrate)) + atexit.register(dev.quit) + print(dev.identify()) + if args.stream != 0: + print(f"Streaming snapshots for {args.element} (Ctrl-C to stop)...") + try: + stream = ( + dev.measure_harmonic_snapshots(args.element, max_order=min(args.max_order, 10)) + if args.harmonics + else dev.measure_snapshots(args.element) + ) + for count, snapshot in enumerate(stream, 1): + print(f"[{count}]", snapshot) + if args.stream > 0 and count >= args.stream: + break + except KeyboardInterrupt: + pass + else: + if args.harmonics: + print( + "Snapshot with harmonics:", + dev.measure_harmonic_snapshot(args.element, max_order=min(args.max_order, 10)), + ) + else: + print("Snapshot (V/A/degrees/W/var):", dev.measure_snapshot(args.element)) + if args.element in ("1A", "1B", "1C") and not args.harmonics: + print("Harmonics (orders 1..max_order, V/A):", dev.measure_harmonics(args.element, args.max_order)) + os.environ["PYTHONINSPECT"] = "1" diff --git a/pyproject.toml b/pyproject.toml index 7917132..872fd05 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "scpi" -version = "2.6.1" +version = "2.7.0" description = "Transport-independent SCPI command sender/parser and device base classes" authors = [{name = "Eero af Heurlin", email = "eero.afheurlin@iki.fi"}] license = "LGPL-2.1-or-later" diff --git a/src/scpi/__init__.py b/src/scpi/__init__.py index e84fab0..9f8dd91 100644 --- a/src/scpi/__init__.py +++ b/src/scpi/__init__.py @@ -1,7 +1,7 @@ """SCPI module, the scpi class implements the base command set, devices may extend it. transports are separate from devices (so you can use for example hp6632b with either serial port or GPIB)""" -__version__ = "2.6.1" # NOTE Use `uv run --locked bump-my-version bump patch` to bump versions correctly +__version__ = "2.7.0" # NOTE Use `uv run --locked bump-my-version bump patch` to bump versions correctly from .errors import CommandError from .scpi import SCPIDevice, SCPIProtocol diff --git a/src/scpi/devices/owh9830.py b/src/scpi/devices/owh9830.py new file mode 100644 index 0000000..f1252a9 --- /dev/null +++ b/src/scpi/devices/owh9830.py @@ -0,0 +1,380 @@ +"""OWON OWH9830 measurements (programming manual, chapters 4 and 7).""" + +import asyncio +from dataclasses import dataclass, field +from decimal import Decimal, InvalidOperation +from typing import Any, ClassVar, Self, override + +import serial as pyserial + +from ..scpi import COMMAND_DEFAULT_TIMEOUT, SCPIDevice +from ..transports.rs232 import RS232Transport + + +class _SerialTransport(RS232Transport): + """OWON serial replies end with LF.""" + + _terminator: ClassVar[bytes] = b"\n" + + +@dataclass +class OWH9830(SCPIDevice): + """RMS measurements using elements 1A, 1B, 1C and 1sigma. + + Commands are spaced by more than 100 ms as required by the manual. + Generic SYST:ERR? checks are unsupported, including on timeouts. + Use one instance per instrument; AIOWrapper supports blocking calls. + """ + + use_safe_variants: bool = field(default=False) + _io_lock: asyncio.Lock = field(default_factory=asyncio.Lock, init=False, repr=False) + _harmonic_lock: asyncio.Lock = field(default_factory=asyncio.Lock, init=False, repr=False) + _last_command: float = field(default=float("-inf"), init=False, repr=False) + _snapshot_elements: tuple[str, ...] | None = field(default=None, init=False, repr=False) + _snapshot_harmonics: int = field(default=0, init=False, repr=False) + _snapshot_period: float = field(default=0.5, init=False, repr=False) + _snapshot_ready_at: float = field(default=0, init=False, repr=False) + _streaming_snapshots: bool = field(default=False, init=False, repr=False) + + def __post_init__(self) -> None: + super().__post_init__() + if self.use_safe_variants: + raise ValueError("OWH9830 does not support generic SYST:ERR? checks") + + async def _pace(self) -> None: + loop = asyncio.get_running_loop() + await asyncio.sleep(max(0, self._last_command + 0.11 - loop.time())) + self._last_command = loop.time() + + @override + async def ask( + self, command: str, cmd_timeout: float = COMMAND_DEFAULT_TIMEOUT, abort_on_timeout: bool = False + ) -> str: + """Pace queries without issuing unsupported error queries on timeout.""" + async with self._io_lock: + return await self._ask(command, cmd_timeout, abort_on_timeout) + + async def _ask( + self, command: str, cmd_timeout: float = COMMAND_DEFAULT_TIMEOUT, abort_on_timeout: bool = False + ) -> str: + """Query while the caller owns _io_lock.""" + await self._pace() + return await self.protocol.ask(command, cmd_timeout, abort_on_timeout, auto_check_error=False) + + @override + async def command( + self, command: str, cmd_timeout: float = COMMAND_DEFAULT_TIMEOUT, abort_on_timeout: bool = False + ) -> None: + """Pace configuration commands; no serial BREAK is sent by default.""" + async with self._io_lock: + await self._command(command, cmd_timeout, abort_on_timeout) + + async def _command( + self, command: str, cmd_timeout: float = COMMAND_DEFAULT_TIMEOUT, abort_on_timeout: bool = False + ) -> None: + """Configure while the caller owns _io_lock.""" + normalized = command.strip().upper().lstrip(":") + if self._streaming_snapshots and normalized.startswith(("NUM", "RATE", "*RST", "DISP", "HARM", "HOLD")): + raise RuntimeError(f"Cannot execute {command.strip()!r} while snapshot iterator is active") + if normalized.startswith(("NUM", "RATE", "*RST", "DISP", "HARM")): + self._snapshot_elements = None + await self._pace() + await self.protocol.command(command, cmd_timeout, abort_on_timeout, auto_check_error=False) + + @staticmethod + def _element(element: str, *, harmonics: bool = False) -> str: + normalized = element.upper() + if normalized not in ("1A", "1B", "1C", "1SIGMA"): + raise ValueError("element must be 1A, 1B, 1C or 1sigma") + if harmonics and normalized == "1SIGMA": + raise ValueError("OWH9830 harmonics support 1A, 1B and 1C only") + return normalized + + async def _values(self, command: str, count: int = 1) -> tuple[Decimal, ...]: + return self._numbers(await self.ask(command), command, count) + + @staticmethod + def _numbers(response: str, command: str, count: int = 1) -> tuple[Decimal, ...]: + # Firmware V1.2.0 appends a comma to harmonic lists. + fields = response.strip().removesuffix(",").split(",") + try: + values = tuple(Decimal(value.strip()) for value in fields) + except InvalidOperation as exc: + raise ValueError(f"Invalid OWH9830 response to {command!r}: {response!r}; check meter LOCAL state") from exc + if len(values) != count or not all(value.is_finite() for value in values): + raise ValueError(f"Expected {count} finite values for {command!r}, got {response!r}") + return values + + async def measure_voltage(self, element: str = "1A") -> Decimal: + """Return RMS voltage in volts.""" + return (await self._values(f":MEAS:VOLT:ELEMENT{self._element(element)}?"))[0] + + async def measure_current(self, element: str = "1A") -> Decimal: + """Return RMS current in amps.""" + return (await self._values(f":MEAS:CURR:ELEMENT{self._element(element)}?"))[0] + + async def measure_phase_angle(self, element: str = "1A") -> Decimal: + """Return the voltage/current phase angle in degrees.""" + return (await self._values(f":MEAS:PHAS:ELEMENT{self._element(element)}?"))[0] + + async def measure_real_power(self, element: str = "1A") -> Decimal: + """Return real (active) power in watts.""" + return (await self._values(f":MEAS:POW:REAL:ELEMENT{self._element(element)}?"))[0] + + async def measure_reactive_power(self, element: str = "1A") -> Decimal: + """Return reactive power in var.""" + return (await self._values(f":MEAS:POW:REAC:ELEMENT{self._element(element)}?"))[0] + + def _snapshot_selection(self, element: str, max_harmonics: int = 0) -> tuple[str, ...]: + normalized = element.upper() + if normalized == "1A-C": + elements = ("1A", "1B", "1C") + else: + normalized = self._element(element, harmonics=max_harmonics > 0) + elements = ("1sigma" if normalized == "1SIGMA" else normalized,) + if max_harmonics != 0: + if "1sigma" in elements: + raise ValueError("OWH9830 harmonics support 1A, 1B and 1C only") + if isinstance(max_harmonics, bool) or not isinstance(max_harmonics, int) or not 1 <= max_harmonics <= 10: + raise ValueError("max_harmonics must be an integer from 1 to 10") + return elements + + async def measure_snapshot(self, element: str = "1A", max_harmonics: int = 0) -> dict[str, dict[str, Any]]: + """Read voltage/current/phase/real/reactive power (and optional harmonics) for one element or 1A-C. + + Returns an element-keyed dictionary, with units V, A, degrees, W and var. + If max_harmonics > 0, includes 'harmonics' with 'voltage' and 'current' tuples (orders 1..max_harmonics). + Unavailable fields are None. Configures and owns the numeric display page and harmonic setup; + repeated calls with the same selection reuse its configuration. + HOLD freezes reported values for the read and its previous state is restored. + Allows one update period before freezing a new frame. Changing selection + while HOLD is already ON raises ValueError; unchanged selections remain held. + Measurement mode is preserved (or switched to HARMONIC when harmonics are requested). + This does not prove simultaneous ADC sampling. + """ + if self._streaming_snapshots: + raise RuntimeError("Cannot execute measure_snapshot while snapshot iterator is active") + elements = self._snapshot_selection(element, max_harmonics) + async with self._io_lock: + hold = (await self._ask(":HOLD?")).strip().upper() + if hold not in ("ON", "OFF", "1", "0"): + raise ValueError(f"Unexpected OWH9830 HOLD state: {hold!r}") + release_hold = hold in ("OFF", "0") + if self._snapshot_elements != elements or self._snapshot_harmonics != max_harmonics: + if not release_hold: + raise ValueError("Release HOLD before changing snapshot selection") + await self._configure_snapshot(elements, max_harmonics) + try: + if release_hold: + await asyncio.sleep(max(0, self._snapshot_ready_at - asyncio.get_running_loop().time())) + await self._command(":HOLD ON") + return await self._snapshot_values(elements, max_harmonics) + except BaseException: + self._snapshot_elements = None + raise + finally: + if release_hold: + await self._command(":HOLD OFF") + self._snapshot_ready_at = asyncio.get_running_loop().time() + self._snapshot_period + 0.11 + + async def measure_harmonic_snapshot(self, element: str = "1A", max_order: int = 10) -> dict[str, dict[str, Any]]: + """Read snapshot including harmonics up to max_order (1..10) for one element or 1A-C.""" + if isinstance(max_order, bool) or not isinstance(max_order, int) or not 1 <= max_order <= 10: + raise ValueError("max_order must be an integer from 1 to 10") + return await self.measure_snapshot(element, max_harmonics=max_order) + + async def _configure_snapshot(self, elements: tuple[str, ...], max_harmonics: int = 0) -> None: + """Select numeric slots and harmonic configuration while caller owns _io_lock and HOLD is OFF.""" + if max_harmonics > 0: + mode = (await self._ask(":DISP:MOD?")).strip().upper() + if mode in ("NORM", "0"): + await self._command(":DISP:MOD HARMONIC") + elif mode not in ("HARMONIC", "1"): + raise ValueError(f"Unexpected OWH9830 measurement mode: {mode!r}") + for element in elements: + await self._command(f":HARM:ORD:ELEMENT{element} 1,{max_harmonics}") + functions = ("U", "I", "pha", "P", "Q") + channels = {"1A": 1, "1B": 2, "1C": 3, "1sigma": 0} + await self._command(f":NUM:NORM:ITEM {16 if len(elements) == 3 else 8}ITEM") + for index, selected in enumerate(elements): + for offset, function in enumerate(functions, 1): + slot = index * len(functions) + offset + await self._command(f":NUM:NORM:OPTION {slot},{function},{channels[selected]}") + await self._command(f":NUM:NORM:NUM {len(elements) * len(functions)}") + period = float((await self._ask(":RATE?")).strip().removesuffix("s")) + if period not in (0.1, 0.2, 0.5, 1, 2, 5): + raise ValueError(f"Unexpected OWH9830 update period: {period!r}") + self._snapshot_period = period + self._snapshot_ready_at = asyncio.get_running_loop().time() + period + 0.11 + self._snapshot_elements = elements + self._snapshot_harmonics = max_harmonics + + async def _snapshot_values(self, elements: tuple[str, ...], max_harmonics: int = 0) -> dict[str, dict[str, Any]]: + """Read numeric slots and harmonic lists while caller owns _io_lock.""" + count = len(elements) * 5 + command = ":NUM:NORM:VAL?" + response = await self._ask(command) + # V1.2.0 has a 128-byte output buffer, including LF. + if len(response) >= 127: + # Indexed NUM slot reads disagree with bulk data on this firmware. + quantities = ("VOLT", "CURR", "PHAS", "POW:REAL", "POW:REAC") + replies = [ + await self._ask(f":MEAS:{quantity}:ELEMENT{selected.upper()}?") + for selected in elements + for quantity in quantities + ] + else: + replies = response.strip().removesuffix(",").split(",") + if len(replies) != count: + raise ValueError(f"Expected {count} snapshot values, got {response!r}") + scalar_values = [None if value.strip() == "----" else self._numbers(value, command)[0] for value in replies] + + harmonics_data: dict[str, dict[str, tuple[Decimal | None, ...]]] = {} + if max_harmonics > 0: + for element in elements: + v_cmd = f":HARM:LIST:VAL:ELEMENT{element}? VOLT" + i_cmd = f":HARM:LIST:VAL:ELEMENT{element}? CURR" + harmonics_data[element] = { + "voltage": self._numbers(await self._ask(v_cmd), v_cmd, max_harmonics), + "current": self._numbers(await self._ask(i_cmd), i_cmd, max_harmonics), + } + + fields = ("voltage", "current", "phase_angle", "real_power", "reactive_power") + result: dict[str, dict[str, Any]] = {} + for index, selected in enumerate(elements): + data: dict[str, Any] = dict(zip(fields, scalar_values[index * 5 : (index + 1) * 5], strict=True)) + if max_harmonics > 0: + data["harmonics"] = harmonics_data[selected] + result[selected] = data + return result + + async def measure_harmonics(self, element: str = "1A", max_order: int = 7) -> dict[str, tuple[Decimal | None, ...]]: + """Return voltage (V) and current (A) RMS harmonics, orders 1..max_order. + + Index zero is the fundamental. max_order must be an integer 1..63. + Enables harmonic measurement mode and leaves it enabled so data keeps + updating. Scalar measurements also work in this mode. Sigma is unsupported. + Updates the display order range. Reads are sequential, not an atomic snapshot. + Above order 10, unavailable readings (----) are represented by None. + """ + if self._streaming_snapshots: + raise RuntimeError("Cannot measure harmonics while snapshot iterator is active") + element = self._element(element, harmonics=True) + if isinstance(max_order, bool) or not isinstance(max_order, int) or not 1 <= max_order <= 63: + raise ValueError("max_order must be an integer from 1 to 63") + async with self._harmonic_lock: + mode = (await self.ask(":DISP:MOD?")).strip().upper() + if mode in ("NORM", "0"): + await self.command(":DISP:MOD HARMONIC") + elif mode not in ("HARMONIC", "1"): + raise ValueError(f"Unexpected OWH9830 measurement mode: {mode!r}") + await self.command(f":HARM:ORD:ELEMENT{element} 1,{max_order}") + if mode in ("NORM", "0") or max_order > 10: + period = float((await self.ask(":RATE?")).strip().removesuffix("s")) + if period not in (0.1, 0.2, 0.5, 1, 2, 5): + raise ValueError(f"Unexpected OWH9830 update period: {period!r}") + await asyncio.sleep(period + 0.11) + if max_order <= 10: + return { + "voltage": await self._values(f":HARM:LIST:VAL:ELEMENT{element}? VOLT", max_order), + "current": await self._values(f":HARM:LIST:VAL:ELEMENT{element}? CURR", max_order), + } + result: dict[str, list[Decimal | None]] = {"voltage": [], "current": []} + # V1.2.0 list replies truncate at 128 bytes and can omit separators. + # ponytail: two queries per order above 10; batch when firmware lists are reliable. + for order in range(1, max_order + 1): + for name, quantity in (("voltage", "VOLT"), ("current", "CURR")): + command = f":MEAS:{quantity}:HARM:ORDER:ELEMENT{element}? {order}" + response = await self.ask(command) + result[name].append(None if response.strip() == "----" else self._numbers(response, command)[0]) + return {name: tuple(values) for name, values in result.items()} + + def measure_snapshots(self, element: str = "1A", max_harmonics: int = 0) -> "SnapshotStream": + """Stream snapshots as fast as possible using an async iterator. + + Configures numeric display slots once and reuses them on every poll. + If max_harmonics > 0, includes harmonics up to max_harmonics (1..10) on every poll. + Commands that change instrument setup refuse to execute while active. + """ + elements = self._snapshot_selection(element, max_harmonics) + return SnapshotStream(self, elements, max_harmonics) + + stream_snapshots = measure_snapshots + snapshots = measure_snapshots + + def measure_harmonic_snapshots(self, element: str = "1A", max_order: int = 10) -> "SnapshotStream": + """Stream snapshots with harmonics up to max_order (1..10) using an async iterator.""" + if isinstance(max_order, bool) or not isinstance(max_order, int) or not 1 <= max_order <= 10: + raise ValueError("max_order must be an integer from 1 to 10") + return self.measure_snapshots(element, max_harmonics=max_order) + + stream_harmonic_snapshots = measure_harmonic_snapshots + harmonic_snapshots = measure_harmonic_snapshots + + +class SnapshotStream: + """Async iterator for rapid OWH9830 snapshot measurement streaming.""" + + def __init__(self, dev: OWH9830, elements: tuple[str, ...], max_harmonics: int = 0) -> None: + self._dev = dev + self._elements = elements + self._max_harmonics = max_harmonics + self._closed = False + self._configured = False + + def __aiter__(self) -> Self: + return self + + async def __anext__(self) -> dict[str, dict[str, Any]]: + if self._closed: + raise StopAsyncIteration + if not self._configured: + if self._dev._streaming_snapshots: + raise RuntimeError("Snapshot iterator is already active") + try: + async with self._dev._io_lock: + if ( + self._dev._snapshot_elements != self._elements + or self._dev._snapshot_harmonics != self._max_harmonics + ): + hold = (await self._dev._ask(":HOLD?")).strip().upper() + if hold not in ("OFF", "0"): + raise ValueError("Release HOLD before changing snapshot selection") + await self._dev._configure_snapshot(self._elements, self._max_harmonics) + self._dev._streaming_snapshots = True + except BaseException: + self._dev._streaming_snapshots = False + self._closed = True + raise + self._configured = True + + async with self._dev._io_lock: + try: + return await self._dev._snapshot_values(self._elements, self._max_harmonics) + except BaseException: + self._dev._snapshot_elements = None + await self.aclose() + raise + + async def aclose(self) -> None: + if not self._closed: + self._closed = True + if self._configured: + self._dev._streaming_snapshots = False + + async def __aenter__(self) -> Self: + return self + + async def __aexit__(self, exc_type: object, exc_val: object, exc_tb: object) -> None: + await self.aclose() + + def __del__(self) -> None: + if self._configured and not self._closed: + self._dev._streaming_snapshots = False + self._closed = True + + +def serial(serial_url: str, baudrate: int = 115200, **kwargs: Any) -> OWH9830: + """Connect directly by RS-232 (SCPI, 8N1); baudrate must match the meter.""" + port = pyserial.serial_for_url(serial_url.strip(), baudrate=baudrate, **kwargs) + return OWH9830(_SerialTransport(serialdevice=port)) diff --git a/src/scpi/wrapper.py b/src/scpi/wrapper.py index 6d19280..d6bd748 100644 --- a/src/scpi/wrapper.py +++ b/src/scpi/wrapper.py @@ -4,11 +4,50 @@ import functools import inspect import logging -from typing import Any +from collections.abc import Callable +from typing import Any, Self LOGGER = logging.getLogger(__name__) +class _SyncAsyncIterator: + """Synchronous iterator wrapping an asynchronous iterator or generator.""" + + def __init__( + self, + async_iter: Any, + loop: asyncio.AbstractEventLoop, + is_closed: Callable[[], bool], + ) -> None: + self._async_iter = async_iter + self._loop = loop + self._is_closed = is_closed + + def __iter__(self) -> Self: + return self + + def __next__(self) -> Any: + if self._is_closed(): + raise RuntimeError("Wrapper is closed") + try: + return self._loop.run_until_complete(anext(self._async_iter)) + except StopAsyncIteration: + raise StopIteration + + def close(self) -> None: + if hasattr(self._async_iter, "aclose") and not self._loop.is_closed(): + self._loop.run_until_complete(self._async_iter.aclose()) + + def __enter__(self) -> Self: + return self + + def __exit__(self, exc_type: object, exc_val: object, exc_tb: object) -> None: + self.close() + + def __del__(self) -> None: + self.close() + + class AIOWrapper: """Wraps all coroutine methods into asyncio run_until_complete calls""" @@ -33,7 +72,7 @@ def loop(self) -> asyncio.AbstractEventLoop: return self._loop def __getattr__(self, item: str) -> Any: - """Get a memeber, if it's a coroutine autowrap it to eventloop run""" + """Get a member, wrapping coroutines and async iterators for synchronous use.""" orig = getattr(self._device, item) if inspect.iscoroutinefunction(orig): @@ -47,6 +86,19 @@ def wrapped(*args: Any, **kwargs: Any) -> Any: return self._loop.run_until_complete(waitable) return wrapped + if callable(orig): + + @functools.wraps(orig) + def wrapped_call(*args: Any, **kwargs: Any) -> Any: + nonlocal self + if self._closed: + raise RuntimeError("Wrapper is closed") + res = orig(*args, **kwargs) + if hasattr(res, "__aiter__") and hasattr(res, "__anext__"): + return _SyncAsyncIterator(res, self._loop, lambda: self._closed) + return res + + return wrapped_call return orig def __dir__(self) -> Any: diff --git a/tests/test_owh9830.py b/tests/test_owh9830.py new file mode 100644 index 0000000..a6b2836 --- /dev/null +++ b/tests/test_owh9830.py @@ -0,0 +1,425 @@ +"""OWH9830 command mapping, firmware replies and request validation.""" + +import asyncio +from decimal import Decimal +from typing import cast +from unittest.mock import AsyncMock, Mock + +import pytest + +from scpi.devices.owh9830 import OWH9830 +from scpi.scpi import SCPIProtocol +from scpi.transports.baseclass import BaseTransport + + +@pytest.mark.asyncio +async def test_measurements() -> None: + protocol = Mock(spec=SCPIProtocol, transport=Mock(spec=BaseTransport)) + protocol.ask = AsyncMock(return_value=" 1.25\r\n") + protocol.command = AsyncMock() + dev = OWH9830(cast(SCPIProtocol, protocol)) + dev._pace = AsyncMock() + methods = ( + (dev.measure_voltage, "VOLT"), + (dev.measure_current, "CURR"), + (dev.measure_phase_angle, "PHAS"), + (dev.measure_real_power, "POW:REAL"), + (dev.measure_reactive_power, "POW:REAC"), + ) + for method, command in methods: + for element in ("1A", "1B", "1C", "1sigma"): + assert await method(element) == Decimal("1.25") + protocol.ask.assert_awaited_with( + f":MEAS:{command}:ELEMENT{element.upper()}?", 1.0, False, auto_check_error=False + ) + + for bad_element in ("1", "2", "1D", "1A;*RST"): + with pytest.raises(ValueError, match="element must"): + await dev.measure_voltage(bad_element) + for bad_order in (0, 64, True, 2.5): + with pytest.raises(ValueError, match="max_order"): + await dev.measure_harmonics(max_order=cast(int, bad_order)) + with pytest.raises(ValueError, match="harmonics support"): + await dev.measure_harmonics("1sigma") + assert protocol.ask.await_count == 20 + protocol.command.assert_not_awaited() + + for response in ("", "\n", "garbage", "1,,2", "NaN", "Infinity", "1,2"): + protocol.ask.return_value = response + with pytest.raises(ValueError): + await dev.measure_voltage() + + # Two concurrent callers must not overwrite each other's harmonic range. + protocol.ask.side_effect = ["HARMONIC", "230, 0.1, \n", "1.2, 0.001, \n", "HARMONIC", "230", "1.2"] + harmonics, fundamental = await asyncio.gather( + dev.measure_harmonics("1B", max_order=2), dev.measure_harmonics("1C", max_order=1) + ) + assert harmonics == {"voltage": (Decimal("230"), Decimal("0.1")), "current": (Decimal("1.2"), Decimal("0.001"))} + assert fundamental == {"voltage": (Decimal("230"),), "current": (Decimal("1.2"),)} + assert [call.args[0] for call in protocol.command.await_args_list] == [ + ":HARM:ORD:ELEMENT1B 1,2", + ":HARM:ORD:ELEMENT1C 1,1", + ] + assert [call.args[0] for call in protocol.ask.await_args_list[-6:]] == [ + ":DISP:MOD?", + ":HARM:LIST:VAL:ELEMENT1B? VOLT", + ":HARM:LIST:VAL:ELEMENT1B? CURR", + ":DISP:MOD?", + ":HARM:LIST:VAL:ELEMENT1C? VOLT", + ":HARM:LIST:VAL:ELEMENT1C? CURR", + ] + protocol.command.reset_mock() + protocol.ask.side_effect = ["HARMONIC", "0.1s"] + [ + " ----" if order == 25 else str(order) for order in range(1, 64) for _ in range(2) + ] + result = await dev.measure_harmonics(max_order=63) + assert ( + result["voltage"] + == result["current"] + == tuple(None if order == 25 else Decimal(order) for order in range(1, 64)) + ) + assert protocol.ask.await_args_list[-1].args[0] == ":MEAS:CURR:HARM:ORDER:ELEMENT1A? 63" + protocol.command.assert_awaited_once_with(":HARM:ORD:ELEMENT1A 1,63", 1.0, False, auto_check_error=False) + protocol.ask.side_effect = None + protocol.ask.side_effect = ["HARMONIC", "1,2,"] + with pytest.raises(ValueError, match="Expected 3 finite"): + await dev.measure_harmonics(max_order=3) + protocol.ask.side_effect = TimeoutError + with pytest.raises(TimeoutError): + await dev.measure_voltage() + protocol.ask.assert_awaited_with(":MEAS:VOLT:ELEMENT1A?", 1.0, False, auto_check_error=False) + + +@pytest.mark.asyncio +async def test_enable_harmonic_mode() -> None: + protocol = Mock(spec=SCPIProtocol, transport=Mock(spec=BaseTransport)) + protocol.ask = AsyncMock(side_effect=["NORM", "0.1s", "230", "1.2"]) + protocol.command = AsyncMock() + dev = OWH9830(cast(SCPIProtocol, protocol)) + dev._pace = AsyncMock() + result = await dev.measure_harmonics(max_order=1) + assert result == {"voltage": (Decimal("230"),), "current": (Decimal("1.2"),)} + assert [call.args[0] for call in protocol.command.await_args_list] == [ + ":DISP:MOD HARMONIC", + ":HARM:ORD:ELEMENT1A 1,1", + ] + + +@pytest.mark.asyncio +async def test_command_spacing() -> None: + times: list[float] = [] + + async def reply(*args: object, **kwargs: object) -> str: + times.append(asyncio.get_running_loop().time()) + return "1" + + protocol = Mock(spec=SCPIProtocol, transport=Mock(spec=BaseTransport)) + protocol.ask = AsyncMock(side_effect=reply) + dev = OWH9830(cast(SCPIProtocol, protocol)) + await asyncio.gather(dev.measure_voltage(), dev.measure_current()) + assert times[1] - times[0] > 0.1 + + +@pytest.mark.asyncio +async def test_snapshot_mapping_and_cache() -> None: + protocol = Mock(spec=SCPIProtocol, transport=Mock(spec=BaseTransport)) + protocol.ask = AsyncMock(side_effect=["OFF", "0.1s", "221.5,3.03,0.61,671,7", "OFF", "221.4,3.02,0.60,670,6"]) + hold_times: list[float] = [] + + async def command(command: str, *args: object, **kwargs: object) -> None: + if command.startswith(":HOLD "): + hold_times.append(asyncio.get_running_loop().time()) + + protocol.command = AsyncMock(side_effect=command) + dev = OWH9830(cast(SCPIProtocol, protocol)) + dev._pace = AsyncMock() + first, second = await asyncio.gather(dev.measure_snapshot(), dev.measure_snapshot("1a")) + assert first == { + "1A": { + "voltage": Decimal("221.5"), + "current": Decimal("3.03"), + "phase_angle": Decimal("0.61"), + "real_power": Decimal("671"), + "reactive_power": Decimal("7"), + } + } + assert second["1A"]["real_power"] == Decimal("670") + assert hold_times[2] - hold_times[1] > 0.1 # Let a fresh frame update between HOLD cycles. + assert [call.args[0] for call in protocol.command.await_args_list] == [ + ":NUM:NORM:ITEM 8ITEM", + ":NUM:NORM:OPTION 1,U,1", + ":NUM:NORM:OPTION 2,I,1", + ":NUM:NORM:OPTION 3,pha,1", + ":NUM:NORM:OPTION 4,P,1", + ":NUM:NORM:OPTION 5,Q,1", + ":NUM:NORM:NUM 5", + ":HOLD ON", + ":HOLD OFF", + ":HOLD ON", + ":HOLD OFF", + ] + protocol.command.reset_mock() + all_reply = "221.5,3.03,0.61,671,7,0,0,----,----,----,0,0,----,----,----" + protocol.ask.side_effect = ["OFF", "0.1s", all_reply] + all_phases = await dev.measure_snapshot("1a-c") + assert tuple(all_phases) == ("1A", "1B", "1C") + assert all_phases["1B"]["voltage"] == Decimal("0") + assert all_phases["1B"]["phase_angle"] is None + assert all_phases["1C"]["reactive_power"] is None + commands = [call.args[0] for call in protocol.command.await_args_list] + assert commands[0] == ":NUM:NORM:ITEM 16ITEM" + assert commands[6] == ":NUM:NORM:OPTION 6,U,2" + assert commands[11] == ":NUM:NORM:OPTION 11,U,3" + assert commands[-3] == ":NUM:NORM:NUM 15" + protocol.command.reset_mock() + protocol.ask.side_effect = ["ON", all_reply] + assert await dev.measure_snapshot("1A-C") == all_phases + protocol.command.assert_not_awaited() # Existing HOLD ON stays ON. + protocol.ask.side_effect = ["ON"] + with pytest.raises(ValueError, match="Release HOLD"): + await dev.measure_snapshot("1A") + protocol.ask.side_effect = ["OFF", "0.1s", "221,3,0.6,670,0"] + sigma = await dev.measure_snapshot("1SIGMA") + assert tuple(sigma) == ("1sigma",) + assert protocol.command.await_args_list[-4].args[0] == ":NUM:NORM:OPTION 5,Q,0" + await dev.command(":NUM:NORM:OPTION 1,I,1") + assert dev._snapshot_elements is None + for invalid in ("1", "2", "1A-D", "1A;*RST"): + with pytest.raises(ValueError, match="element must"): + await dev.measure_snapshot(invalid) + + +@pytest.mark.asyncio +@pytest.mark.parametrize("failure", ["1,2", "1,2,NaN,4,5", "1,2,,4,5", TimeoutError(), asyncio.CancelledError()]) +async def test_snapshot_releases_hold_on_failure(failure: str | BaseException) -> None: + protocol = Mock(spec=SCPIProtocol, transport=Mock(spec=BaseTransport)) + protocol.ask = AsyncMock(side_effect=["OFF", "0.1s", failure]) + protocol.command = AsyncMock() + dev = OWH9830(cast(SCPIProtocol, protocol)) + dev._pace = AsyncMock() + error = ValueError if isinstance(failure, str) else type(failure) + with pytest.raises(error): + await dev.measure_snapshot() + assert [call.args[0] for call in protocol.command.await_args_list[-2:]] == [":HOLD ON", ":HOLD OFF"] + assert dev._snapshot_elements is None + + +@pytest.mark.asyncio +async def test_snapshot_long_reply_and_io_exclusion() -> None: + commands: list[str] = [] + replies = iter(["OFF", "0.1s", "1,2,3,4," + "0" * 119, "11", "22", "----", "44", "55", "230"]) + + async def ask(command: str, *args: object, **kwargs: object) -> str: + commands.append(command) + await asyncio.sleep(0) + return next(replies) + + async def command(command: str, *args: object, **kwargs: object) -> None: + commands.append(command) + await asyncio.sleep(0) + + protocol = Mock(spec=SCPIProtocol, transport=Mock(spec=BaseTransport)) + protocol.ask = AsyncMock(side_effect=ask) + protocol.command = AsyncMock(side_effect=command) + dev = OWH9830(cast(SCPIProtocol, protocol)) + dev._pace = AsyncMock() + snapshot, voltage = await asyncio.gather(dev.measure_snapshot(), dev.measure_voltage()) + assert snapshot["1A"]["voltage"] == Decimal("11") + assert snapshot["1A"]["phase_angle"] is None + assert snapshot["1A"]["reactive_power"] == Decimal("55") + assert voltage == Decimal("230") + assert commands[-7:] == [ + ":MEAS:VOLT:ELEMENT1A?", + ":MEAS:CURR:ELEMENT1A?", + ":MEAS:PHAS:ELEMENT1A?", + ":MEAS:POW:REAL:ELEMENT1A?", + ":MEAS:POW:REAC:ELEMENT1A?", + ":HOLD OFF", + ":MEAS:VOLT:ELEMENT1A?", + ] + + +@pytest.mark.asyncio +async def test_measure_snapshots_stream() -> None: + protocol = Mock(spec=SCPIProtocol, transport=Mock(spec=BaseTransport)) + protocol.ask = AsyncMock( + side_effect=[ + "OFF", + "0.1s", + "220.0,3.00,0.60,660,10", + "220.1,3.01,0.60,661,11", + "220.2,3.02,0.60,662,12", + ] + ) + protocol.command = AsyncMock() + dev = OWH9830(cast(SCPIProtocol, protocol)) + dev._pace = AsyncMock() + + # 1. Single-element streaming setup once and rapid reads + stream = dev.measure_snapshots("1A") + items = [] + async for s in stream: + items.append(s) + if len(items) == 3: + break + + assert len(items) == 3 + assert items[0]["1A"]["voltage"] == Decimal("220.0") + assert items[1]["1A"]["voltage"] == Decimal("220.1") + assert items[2]["1A"]["voltage"] == Decimal("220.2") + + # Slot setup was called once only, then :NUM:NORM:VAL? queries + assert [call.args[0] for call in protocol.command.await_args_list] == [ + ":NUM:NORM:ITEM 8ITEM", + ":NUM:NORM:OPTION 1,U,1", + ":NUM:NORM:OPTION 2,I,1", + ":NUM:NORM:OPTION 3,pha,1", + ":NUM:NORM:OPTION 4,P,1", + ":NUM:NORM:OPTION 5,Q,1", + ":NUM:NORM:NUM 5", + ] + assert [call.args[0] for call in protocol.ask.await_args_list] == [ + ":HOLD?", + ":RATE?", + ":NUM:NORM:VAL?", + ":NUM:NORM:VAL?", + ":NUM:NORM:VAL?", + ] + + # Stream is still active, verify setup-changing commands are blocked + with pytest.raises(RuntimeError, match="Cannot measure harmonics"): + await dev.measure_harmonics() + with pytest.raises(RuntimeError, match="Cannot execute measure_snapshot"): + await dev.measure_snapshot() + with pytest.raises(RuntimeError, match="while snapshot iterator is active"): + await dev.command(":RATE 1s") + with pytest.raises(RuntimeError, match="while snapshot iterator is active"): + await dev.command(":NUM:NORM:NUM 10") + with pytest.raises(RuntimeError, match="while snapshot iterator is active"): + await dev.command("*RST") + with pytest.raises(RuntimeError, match="while snapshot iterator is active"): + await dev.command(":DISP:MOD NORM") + with pytest.raises(RuntimeError, match="while snapshot iterator is active"): + await dev.command(":HOLD ON") + with pytest.raises(RuntimeError, match="while snapshot iterator is active"): + await dev.command(":HARM:ORD:ELEMENT1A 1,7") + + # Concurrent iterator should be refused + stream2 = dev.measure_snapshots("1B") + with pytest.raises(RuntimeError, match="Snapshot iterator is already active"): + await anext(stream2) + + # Non-setup modifying queries work + protocol.ask.side_effect = ["220.5"] + assert await dev.measure_voltage("1A") == Decimal("220.5") + + # Close the stream + await stream.aclose() + + # After close, setup commands and measure_snapshot work again + protocol.ask.side_effect = ["OFF", "220.0,3.00,0.60,660,10"] + protocol.command.reset_mock() + snap = await dev.measure_snapshot("1A") + assert snap["1A"]["voltage"] == Decimal("220.0") + + # Context manager usage + protocol.ask.side_effect = ["220.0,3.00,0.60,660,10"] + async with dev.measure_snapshots("1A") as stream3: + item = await anext(stream3) + assert item["1A"]["voltage"] == Decimal("220.0") + assert dev._streaming_snapshots is True + assert dev._streaming_snapshots is False + + +@pytest.mark.asyncio +async def test_harmonic_snapshot_mapping_and_errors() -> None: + protocol = Mock(spec=SCPIProtocol, transport=Mock(spec=BaseTransport)) + protocol.ask = AsyncMock( + side_effect=[ + "OFF", + "NORM", + "0.1s", + "220.0,3.00,0.60,660,10", + "220.0, 0.1, 0.2, 0.3, 0.4, 0.5, 0.6, 0.7, 0.8, 0.9,", + "3.0, 0.01, 0.02, 0.03, 0.04, 0.05, 0.06, 0.07, 0.08, 0.09,", + ] + ) + protocol.command = AsyncMock() + dev = OWH9830(cast(SCPIProtocol, protocol)) + dev._pace = AsyncMock() + + snap = await dev.measure_harmonic_snapshot("1A", max_order=10) + assert snap["1A"]["voltage"] == Decimal("220.0") + assert snap["1A"]["real_power"] == Decimal("660") + assert snap["1A"]["harmonics"]["voltage"] == tuple( + Decimal(f"0.{i}") if i > 0 else Decimal("220.0") for i in range(10) + ) + assert snap["1A"]["harmonics"]["current"] == tuple( + Decimal(f"0.0{i}") if i > 0 else Decimal("3.0") for i in range(10) + ) + + assert [call.args[0] for call in protocol.command.await_args_list] == [ + ":DISP:MOD HARMONIC", + ":HARM:ORD:ELEMENT1A 1,10", + ":NUM:NORM:ITEM 8ITEM", + ":NUM:NORM:OPTION 1,U,1", + ":NUM:NORM:OPTION 2,I,1", + ":NUM:NORM:OPTION 3,pha,1", + ":NUM:NORM:OPTION 4,P,1", + ":NUM:NORM:OPTION 5,Q,1", + ":NUM:NORM:NUM 5", + ":HOLD ON", + ":HOLD OFF", + ] + + # Error conditions + for bad_order in (0, 11, True, "10", 2.5): + with pytest.raises(ValueError, match="max_order must be an integer from 1 to 10"): + await dev.measure_harmonic_snapshot("1A", max_order=cast(int, bad_order)) + + with pytest.raises(ValueError, match="harmonics support 1A, 1B and 1C only"): + await dev.measure_harmonic_snapshot("1sigma") + + +@pytest.mark.asyncio +async def test_harmonic_snapshots_stream() -> None: + protocol = Mock(spec=SCPIProtocol, transport=Mock(spec=BaseTransport)) + protocol.ask = AsyncMock( + side_effect=[ + "OFF", + "HARMONIC", + "0.1s", + "220.0,3.00,0.60,660,10", + "220.0, 0.1, 0.2, 0.3, 0.4, 0.5, 0.6,", + "3.0, 0.01, 0.02, 0.03, 0.04, 0.05, 0.06,", + "220.1,3.01,0.60,661,11", + "220.1, 0.1, 0.2, 0.3, 0.4, 0.5, 0.6,", + "3.01, 0.01, 0.02, 0.03, 0.04, 0.05, 0.06,", + ] + ) + protocol.command = AsyncMock() + dev = OWH9830(cast(SCPIProtocol, protocol)) + dev._pace = AsyncMock() + + items = [] + async for s in dev.measure_harmonic_snapshots("1A", max_order=7): + items.append(s) + if len(items) == 2: + break + + assert len(items) == 2 + assert items[0]["1A"]["voltage"] == Decimal("220.0") + assert len(items[0]["1A"]["harmonics"]["voltage"]) == 7 + assert items[1]["1A"]["voltage"] == Decimal("220.1") + assert len(items[1]["1A"]["harmonics"]["voltage"]) == 7 + + assert [call.args[0] for call in protocol.command.await_args_list] == [ + ":HARM:ORD:ELEMENT1A 1,7", + ":NUM:NORM:ITEM 8ITEM", + ":NUM:NORM:OPTION 1,U,1", + ":NUM:NORM:OPTION 2,I,1", + ":NUM:NORM:OPTION 3,pha,1", + ":NUM:NORM:OPTION 4,P,1", + ":NUM:NORM:OPTION 5,Q,1", + ":NUM:NORM:NUM 5", + ] diff --git a/tests/test_scpi.py b/tests/test_scpi.py index 2ec5f61..d510d8f 100644 --- a/tests/test_scpi.py +++ b/tests/test_scpi.py @@ -5,4 +5,4 @@ def test_version() -> None: """Make sure version matches expected""" - assert __version__ == "2.6.1" + assert __version__ == "2.7.0" diff --git a/tests/test_serial.py b/tests/test_serial.py index af33c6a..ef26049 100644 --- a/tests/test_serial.py +++ b/tests/test_serial.py @@ -8,6 +8,7 @@ import pytest import serial +from scpi.devices.owh9830 import _SerialTransport from scpi.transports.gpib.prologix import PrologixGPIBTransport from scpi.transports.rs232 import RS232Transport, get from scpi.wrapper import AIOWrapper @@ -43,7 +44,7 @@ async def read_peer(master: int, expected: bytes) -> None: @pytest.mark.asyncio @pytest.mark.parametrize( ("transport_type", "terminator", "startup"), - [(RS232Transport, b"\r\n", b""), (PrologixGPIBTransport, b"\n", STARTUP)], + [(RS232Transport, b"\r\n", b""), (PrologixGPIBTransport, b"\n", STARTUP), (_SerialTransport, b"\n", b"")], ) async def test_serial_exchange(transport_type: type[RS232Transport], terminator: bytes, startup: bytes) -> None: """Verify initialization, framing, fragmented input and replies received before reads.""" diff --git a/tests/test_wrapper.py b/tests/test_wrapper.py index fae2bf3..888b98a 100644 --- a/tests/test_wrapper.py +++ b/tests/test_wrapper.py @@ -65,3 +65,50 @@ def test_shared_loop() -> None: assert controller.identify() == "instrument" finally: controller.quit() + + +def test_wrapper_async_iterator() -> None: + """Async iterators can be used synchronously via the wrapper.""" + + class StreamDevice(Device): + def stream(self): + class AsyncIter: + def __init__(self, dev: Device) -> None: + self.dev = dev + self.count = 0 + self.closed = False + + def __aiter__(self): + return self + + async def __anext__(self): + if self.closed or self.count >= 3: + raise StopAsyncIteration + self.dev.loops.append(asyncio.get_running_loop()) + self.count += 1 + return self.count + + async def aclose(self): + self.closed = True + + async def __aenter__(self): + return self + + async def __aexit__(self, *args: object): + await self.aclose() + + return AsyncIter(self) + + device = StreamDevice() + wrapper = AIOWrapper(device) + try: + results = [] + for item in wrapper.stream(): + results.append(item) + assert results == [1, 2, 3] + + # Context manager and close + with wrapper.stream() as it: + assert next(it) == 1 + finally: + wrapper.quit() diff --git a/uv.lock b/uv.lock index 552f9fb..7ac5c06 100644 --- a/uv.lock +++ b/uv.lock @@ -863,7 +863,7 @@ wheels = [ [[package]] name = "scpi" -version = "2.6.1" +version = "2.7.0" source = { editable = "." } dependencies = [ { name = "pyserial" },