From 6a465b3486e64b317b0b9f0fa47c512466b5f6c9 Mon Sep 17 00:00:00 2001 From: Eero af Heurlin Date: Fri, 2 Oct 2026 13:13:46 +0300 Subject: [PATCH 01/12] chore: ignore local user manuals and handoff files --- .gitignore | 5 +++++ 1 file changed, 5 insertions(+) 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 From 2b44f7e973bfd0095c99832a3c2917866220e037 Mon Sep 17 00:00:00 2001 From: Eero af Heurlin Date: Fri, 2 Oct 2026 13:43:10 +0300 Subject: [PATCH 02/12] feat(owh9830): add power and harmonic readings --- src/scpi/devices/owh9830.py | 149 ++++++++++++++++++++++++++++++++++++ tests/test_owh9830.py | 120 +++++++++++++++++++++++++++++ tests/test_serial.py | 3 +- 3 files changed, 271 insertions(+), 1 deletion(-) create mode 100644 src/scpi/devices/owh9830.py create mode 100644 tests/test_owh9830.py diff --git a/src/scpi/devices/owh9830.py b/src/scpi/devices/owh9830.py new file mode 100644 index 0000000..f2e9d0c --- /dev/null +++ b/src/scpi/devices/owh9830.py @@ -0,0 +1,149 @@ +"""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, 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) + + 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: + 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._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] + + 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. + """ + 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 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/tests/test_owh9830.py b/tests/test_owh9830.py new file mode 100644 index 0000000..cbb7baf --- /dev/null +++ b/tests/test_owh9830.py @@ -0,0 +1,120 @@ +"""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 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.""" From e2bce3ca693412fe0dd7133ecea9d545d5b0ff72 Mon Sep 17 00:00:00 2001 From: Eero af Heurlin Date: Fri, 2 Oct 2026 13:44:03 +0300 Subject: [PATCH 03/12] docs(owh9830): add serial example and usage --- examples/owh9830.rst | 40 ++++++++++++++++++++++++++++++++++++++ examples/owh9830_serial.py | 28 ++++++++++++++++++++++++++ 2 files changed, 68 insertions(+) create mode 100644 examples/owh9830.rst create mode 100644 examples/owh9830_serial.py diff --git a/examples/owh9830.rst b/examples/owh9830.rst new file mode 100644 index 0000000..9864ee7 --- /dev/null +++ b/examples/owh9830.rst @@ -0,0 +1,40 @@ +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 diff --git a/examples/owh9830_serial.py b/examples/owh9830_serial.py new file mode 100644 index 0000000..369e434 --- /dev/null +++ b/examples/owh9830_serial.py @@ -0,0 +1,28 @@ +#!/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"), default="1A") + parser.add_argument("--max-order", type=int, choices=range(1, 64), default=7) + args = parser.parse_args() + dev = AIOWrapper(serial(args.port, baudrate=args.baudrate)) + atexit.register(dev.quit) + print(dev.identify()) + print("Voltage (V):", dev.measure_voltage(args.element)) + print("Current (A):", dev.measure_current(args.element)) + print("Phase angle (degrees):", dev.measure_phase_angle(args.element)) + print("Real power (W):", dev.measure_real_power(args.element)) + print("Reactive power (var):", dev.measure_reactive_power(args.element)) + if args.element != "1sigma": + print("Harmonics (orders 1..max_order, V/A):", dev.measure_harmonics(args.element, args.max_order)) + os.environ["PYTHONINSPECT"] = "1" From 799d5d8531a4f6dc2f5cf0780cc00d9181077204 Mon Sep 17 00:00:00 2001 From: Eero af Heurlin Date: Fri, 2 Oct 2026 14:23:12 +0300 Subject: [PATCH 04/12] feat(owh9830): read held measurement sets --- src/scpi/devices/owh9830.py | 106 ++++++++++++++++++++++++++++++-- tests/test_owh9830.py | 119 ++++++++++++++++++++++++++++++++++++ 2 files changed, 221 insertions(+), 4 deletions(-) diff --git a/src/scpi/devices/owh9830.py b/src/scpi/devices/owh9830.py index f2e9d0c..daa1802 100644 --- a/src/scpi/devices/owh9830.py +++ b/src/scpi/devices/owh9830.py @@ -30,6 +30,9 @@ class OWH9830(SCPIDevice): _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_period: float = field(default=0.5, init=False, repr=False) + _snapshot_ready_at: float = field(default=0, init=False, repr=False) def __post_init__(self) -> None: super().__post_init__() @@ -47,8 +50,14 @@ async def ask( ) -> str: """Pace queries without issuing unsupported error queries on timeout.""" async with self._io_lock: - await self._pace() - return await self.protocol.ask(command, cmd_timeout, abort_on_timeout, auto_check_error=False) + 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( @@ -56,8 +65,16 @@ async def command( ) -> None: """Pace configuration commands; no serial BREAK is sent by default.""" async with self._io_lock: - await self._pace() - await self.protocol.command(command, cmd_timeout, abort_on_timeout, auto_check_error=False) + 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.""" + if command.strip().upper().lstrip(":").startswith(("NUM", "RATE", "*RST")): + 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: @@ -103,6 +120,87 @@ 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] + async def measure_snapshot(self, element: str = "1A") -> dict[str, dict[str, Decimal | None]]: + """Read voltage/current/phase/real/reactive power for one element or 1A-C. + + Returns an element-keyed dictionary, with units V, A, degrees, W and var. + Unavailable fields are None. Configures and owns the numeric display page; + 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. This does not prove simultaneous ADC sampling. + """ + normalized = element.upper() + if normalized == "1A-C": + elements = ("1A", "1B", "1C") + else: + normalized = self._element(element) + elements = ("1sigma" if normalized == "1SIGMA" else normalized,) + fields = ("voltage", "current", "phase_angle", "real_power", "reactive_power") + 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: + if not release_hold: + raise ValueError("Release HOLD before changing snapshot selection") + await self._configure_snapshot(elements) + try: + if release_hold: + await asyncio.sleep(max(0, self._snapshot_ready_at - asyncio.get_running_loop().time())) + await self._command(":HOLD ON") + values = await self._snapshot_values(elements) + 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 + return { + selected: dict(zip(fields, values[index * len(fields) : (index + 1) * len(fields)], strict=True)) + for index, selected in enumerate(elements) + } + + async def _configure_snapshot(self, elements: tuple[str, ...]) -> None: + """Select numeric slots while the caller owns _io_lock and HOLD is OFF.""" + 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 + + async def _snapshot_values(self, elements: tuple[str, ...]) -> list[Decimal | None]: + """Read numeric slots while the caller owns _io_lock and HOLD is active.""" + 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}") + return [None if value.strip() == "----" else self._numbers(value, command)[0] for value in replies] + 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. diff --git a/tests/test_owh9830.py b/tests/test_owh9830.py index cbb7baf..df53deb 100644 --- a/tests/test_owh9830.py +++ b/tests/test_owh9830.py @@ -118,3 +118,122 @@ async def reply(*args: object, **kwargs: object) -> str: 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?", + ] From fea9f5dc8ee0726c97b42c4235882198619193c2 Mon Sep 17 00:00:00 2001 From: Eero af Heurlin Date: Fri, 2 Oct 2026 14:23:33 +0300 Subject: [PATCH 05/12] docs(owh9830): explain snapshot polling --- examples/owh9830.rst | 85 ++++++++++++++++++++++++++++++++++++++ examples/owh9830_serial.py | 10 ++--- 2 files changed, 88 insertions(+), 7 deletions(-) diff --git a/examples/owh9830.rst b/examples/owh9830.rst index 9864ee7..e9a3d51 100644 --- a/examples/owh9830.rst +++ b/examples/owh9830.rst @@ -38,3 +38,88 @@ 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"]) + +The result always maps element names to dictionaries containing ``voltage`` (V), +``current`` (A), ``phase_angle`` (degrees), ``real_power`` (W) and +``reactive_power`` (var). Single ``1A``, ``1B``, ``1C`` and ``1sigma`` selections +are supported, as is the combined ``1A-C`` selection. 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), then reuses it for repeated polls. It queries the +current HOLD state, enables HOLD if necessary, reads the configured slots, 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. + +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 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 example is:: + + uv run --locked python examples/owh9830_serial.py /dev/tty.usbserial-A92TRYDJ --element 1A-C + +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_serial.py b/examples/owh9830_serial.py index 369e434..1ff2653 100644 --- a/examples/owh9830_serial.py +++ b/examples/owh9830_serial.py @@ -12,17 +12,13 @@ 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"), default="1A") + 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) args = parser.parse_args() dev = AIOWrapper(serial(args.port, baudrate=args.baudrate)) atexit.register(dev.quit) print(dev.identify()) - print("Voltage (V):", dev.measure_voltage(args.element)) - print("Current (A):", dev.measure_current(args.element)) - print("Phase angle (degrees):", dev.measure_phase_angle(args.element)) - print("Real power (W):", dev.measure_real_power(args.element)) - print("Reactive power (var):", dev.measure_reactive_power(args.element)) - if args.element != "1sigma": + print("Snapshot (V/A/degrees/W/var):", dev.measure_snapshot(args.element)) + if args.element in ("1A", "1B", "1C"): print("Harmonics (orders 1..max_order, V/A):", dev.measure_harmonics(args.element, args.max_order)) os.environ["PYTHONINSPECT"] = "1" From aeb33229e1ba1241ce81de41b3d5610f0256cbb4 Mon Sep 17 00:00:00 2001 From: Eero af Heurlin Date: Fri, 2 Oct 2026 14:39:00 +0300 Subject: [PATCH 06/12] feat(owh9830): stream fast snapshot measurements via async iterator Co-authored-by: Junie --- src/scpi/devices/owh9830.py | 95 ++++++++++++++++++++++++++++++++++++- tests/test_owh9830.py | 92 +++++++++++++++++++++++++++++++++++ 2 files changed, 185 insertions(+), 2 deletions(-) diff --git a/src/scpi/devices/owh9830.py b/src/scpi/devices/owh9830.py index daa1802..ae8b50b 100644 --- a/src/scpi/devices/owh9830.py +++ b/src/scpi/devices/owh9830.py @@ -3,7 +3,7 @@ import asyncio from dataclasses import dataclass, field from decimal import Decimal, InvalidOperation -from typing import Any, ClassVar, override +from typing import Any, ClassVar, Self, override import serial as pyserial @@ -33,6 +33,7 @@ class OWH9830(SCPIDevice): _snapshot_elements: tuple[str, ...] | None = field(default=None, 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__() @@ -71,7 +72,10 @@ 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.""" - if command.strip().upper().lstrip(":").startswith(("NUM", "RATE", "*RST")): + 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")): self._snapshot_elements = None await self._pace() await self.protocol.command(command, cmd_timeout, abort_on_timeout, auto_check_error=False) @@ -131,6 +135,8 @@ async def measure_snapshot(self, element: str = "1A") -> dict[str, dict[str, Dec while HOLD is already ON raises ValueError; unchanged selections remain held. Measurement mode is preserved. This does not prove simultaneous ADC sampling. """ + if self._streaming_snapshots: + raise RuntimeError("Cannot execute measure_snapshot while snapshot iterator is active") normalized = element.upper() if normalized == "1A-C": elements = ("1A", "1B", "1C") @@ -210,6 +216,8 @@ async def measure_harmonics(self, element: str = "1A", max_order: int = 7) -> di 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") @@ -240,6 +248,89 @@ async def measure_harmonics(self, element: str = "1A", max_order: int = 7) -> di 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") -> "SnapshotStream": + """Stream snapshots as fast as possible using an async iterator. + + Configures numeric display slots once and reuses them on every poll. + Commands that change instrument setup refuse to execute while active. + """ + normalized = element.upper() + if normalized == "1A-C": + elements = ("1A", "1B", "1C") + else: + normalized = self._element(element) + elements = ("1sigma" if normalized == "1SIGMA" else normalized,) + fields = ("voltage", "current", "phase_angle", "real_power", "reactive_power") + return SnapshotStream(self, elements, fields) + + stream_snapshots = measure_snapshots + snapshots = measure_snapshots + + +class SnapshotStream: + """Async iterator for rapid OWH9830 snapshot measurement streaming.""" + + def __init__(self, dev: OWH9830, elements: tuple[str, ...], fields: tuple[str, ...]) -> None: + self._dev = dev + self._elements = elements + self._fields = fields + self._closed = False + self._configured = False + + def __aiter__(self) -> Self: + return self + + async def __anext__(self) -> dict[str, dict[str, Decimal | None]]: + 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: + 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._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: + values = await self._dev._snapshot_values(self._elements) + except BaseException: + self._dev._snapshot_elements = None + await self.aclose() + raise + return { + selected: dict( + zip(self._fields, values[index * len(self._fields) : (index + 1) * len(self._fields)], strict=True) + ) + for index, selected in enumerate(self._elements) + } + + 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.""" diff --git a/tests/test_owh9830.py b/tests/test_owh9830.py index df53deb..f182bc7 100644 --- a/tests/test_owh9830.py +++ b/tests/test_owh9830.py @@ -237,3 +237,95 @@ async def command(command: str, *args: object, **kwargs: object) -> None: ":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 From fea1a732d678c33fce117d41679e3d80ffc7417b Mon Sep 17 00:00:00 2001 From: Eero af Heurlin Date: Fri, 2 Oct 2026 14:39:04 +0300 Subject: [PATCH 07/12] feat(wrapper): support async iterators in blocking wrapper Co-authored-by: Junie --- src/scpi/wrapper.py | 56 +++++++++++++++++++++++++++++++++++++++++-- tests/test_wrapper.py | 47 ++++++++++++++++++++++++++++++++++++ 2 files changed, 101 insertions(+), 2 deletions(-) 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_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() From bab851a1406ac3352bfc51368c6591fd3767a994 Mon Sep 17 00:00:00 2001 From: Eero af Heurlin Date: Fri, 2 Oct 2026 14:39:08 +0300 Subject: [PATCH 08/12] docs(owh9830): document streaming snapshots and add serial example flag Co-authored-by: Junie --- examples/owh9830.rst | 25 +++++++++++++++++++++++++ examples/owh9830_serial.py | 19 ++++++++++++++++--- 2 files changed, 41 insertions(+), 3 deletions(-) diff --git a/examples/owh9830.rst b/examples/owh9830.rst index e9a3d51..5e0c0e2 100644 --- a/examples/owh9830.rst +++ b/examples/owh9830.rst @@ -80,6 +80,31 @@ The interactive all-phase example is:: uv run --locked python examples/owh9830_serial.py /dev/tty.usbserial-A92TRYDJ --element 1A-C +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"]) + +Synchronous code using ``AIOWrapper`` supports streaming using a standard ``for`` loop:: + + for snapshot in dev.measure_snapshots("1A"): + print(snapshot["1A"]) + +The iterator configures the numeric display page slots 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 streaming example is:: + + uv run --locked python examples/owh9830_serial.py /dev/tty.usbserial-A92TRYDJ --stream 5 + Batching experiments -------------------- diff --git a/examples/owh9830_serial.py b/examples/owh9830_serial.py index 1ff2653..9b4f5aa 100644 --- a/examples/owh9830_serial.py +++ b/examples/owh9830_serial.py @@ -14,11 +14,24 @@ 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( + "--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()) - print("Snapshot (V/A/degrees/W/var):", dev.measure_snapshot(args.element)) - if args.element in ("1A", "1B", "1C"): - print("Harmonics (orders 1..max_order, V/A):", dev.measure_harmonics(args.element, args.max_order)) + if args.stream != 0: + print(f"Streaming snapshots for {args.element} (Ctrl-C to stop)...") + try: + for count, snapshot in enumerate(dev.measure_snapshots(args.element), 1): + print(f"[{count}]", snapshot) + if args.stream > 0 and count >= args.stream: + break + except KeyboardInterrupt: + pass + else: + print("Snapshot (V/A/degrees/W/var):", dev.measure_snapshot(args.element)) + if args.element in ("1A", "1B", "1C"): + print("Harmonics (orders 1..max_order, V/A):", dev.measure_harmonics(args.element, args.max_order)) os.environ["PYTHONINSPECT"] = "1" From 113e2d70b1ea425eb424a861cb6b26d7d4431ae2 Mon Sep 17 00:00:00 2001 From: Eero af Heurlin Date: Fri, 2 Oct 2026 14:42:13 +0300 Subject: [PATCH 09/12] docs(owh9830): add async snapshot streaming example Co-authored-by: Junie --- examples/owh9830.rst | 3 ++- examples/owh9830_async_stream.py | 45 ++++++++++++++++++++++++++++++++ 2 files changed, 47 insertions(+), 1 deletion(-) create mode 100644 examples/owh9830_async_stream.py diff --git a/examples/owh9830.rst b/examples/owh9830.rst index 5e0c0e2..ff6d41e 100644 --- a/examples/owh9830.rst +++ b/examples/owh9830.rst @@ -101,9 +101,10 @@ harmonic measurements, single snapshots, or ``:NUM``, ``:RATE``, ``*RST``, execute and raise ``RuntimeError`` to prevent configuration conflicts. Read-only queries (e.g. ``measure_voltage`` or ``ask``) execute safely between polls. -The interactive streaming example is:: +The interactive examples are:: uv run --locked python examples/owh9830_serial.py /dev/tty.usbserial-A92TRYDJ --stream 5 + uv run --locked python examples/owh9830_async_stream.py /dev/tty.usbserial-A92TRYDJ --count 5 Batching experiments -------------------- diff --git a/examples/owh9830_async_stream.py b/examples/owh9830_async_stream.py new file mode 100644 index 0000000..a103290 --- /dev/null +++ b/examples/owh9830_async_stream.py @@ -0,0 +1,45 @@ +#!/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( + "--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} (Ctrl-C to stop)...") + + received = 0 + async for snapshot in dev.measure_snapshots(args.element): + 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) From 50ffdc8f998b9329f2034cd3d9e3aefe5dd03f8b Mon Sep 17 00:00:00 2001 From: Eero af Heurlin Date: Fri, 2 Oct 2026 15:01:05 +0300 Subject: [PATCH 10/12] chore: bump version --- .bumpversion.toml | 2 +- pyproject.toml | 2 +- src/scpi/__init__.py | 2 +- tests/test_scpi.py | 2 +- uv.lock | 2 +- 5 files changed, 5 insertions(+), 5 deletions(-) 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/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/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/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" }, From b983569d84d6170aa395b1da438574fe22dead6c Mon Sep 17 00:00:00 2001 From: Eero af Heurlin Date: Fri, 2 Oct 2026 15:17:00 +0300 Subject: [PATCH 11/12] feat(owh9830): add harmonic snapshots and streaming iterator Co-authored-by: Junie --- src/scpi/devices/owh9830.py | 132 ++++++++++++++++++++++++------------ tests/test_owh9830.py | 94 +++++++++++++++++++++++++ 2 files changed, 181 insertions(+), 45 deletions(-) diff --git a/src/scpi/devices/owh9830.py b/src/scpi/devices/owh9830.py index ae8b50b..f1252a9 100644 --- a/src/scpi/devices/owh9830.py +++ b/src/scpi/devices/owh9830.py @@ -31,6 +31,7 @@ class OWH9830(SCPIDevice): _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) @@ -75,7 +76,7 @@ async def _command( 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")): + 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) @@ -124,40 +125,50 @@ 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] - async def measure_snapshot(self, element: str = "1A") -> dict[str, dict[str, Decimal | None]]: - """Read voltage/current/phase/real/reactive power for one element or 1A-C. + 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. - Unavailable fields are None. Configures and owns the numeric display page; + 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. This does not prove simultaneous ADC sampling. + 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") - normalized = element.upper() - if normalized == "1A-C": - elements = ("1A", "1B", "1C") - else: - normalized = self._element(element) - elements = ("1sigma" if normalized == "1SIGMA" else normalized,) - fields = ("voltage", "current", "phase_angle", "real_power", "reactive_power") + 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: + 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) + 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") - values = await self._snapshot_values(elements) + return await self._snapshot_values(elements, max_harmonics) except BaseException: self._snapshot_elements = None raise @@ -165,13 +176,23 @@ async def measure_snapshot(self, element: str = "1A") -> dict[str, dict[str, Dec if release_hold: await self._command(":HOLD OFF") self._snapshot_ready_at = asyncio.get_running_loop().time() + self._snapshot_period + 0.11 - return { - selected: dict(zip(fields, values[index * len(fields) : (index + 1) * len(fields)], strict=True)) - for index, selected in enumerate(elements) - } - async def _configure_snapshot(self, elements: tuple[str, ...]) -> None: - """Select numeric slots while the caller owns _io_lock and HOLD is OFF.""" + 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") @@ -186,9 +207,10 @@ async def _configure_snapshot(self, elements: tuple[str, ...]) -> None: 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, ...]) -> list[Decimal | None]: - """Read numeric slots while the caller owns _io_lock and HOLD is active.""" + 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) @@ -205,7 +227,26 @@ async def _snapshot_values(self, elements: tuple[str, ...]) -> list[Decimal | No replies = response.strip().removesuffix(",").split(",") if len(replies) != count: raise ValueError(f"Expected {count} snapshot values, got {response!r}") - return [None if value.strip() == "----" else self._numbers(value, command)[0] for value in replies] + 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. @@ -248,39 +289,43 @@ async def measure_harmonics(self, element: str = "1A", max_order: int = 7) -> di 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") -> "SnapshotStream": + 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. """ - normalized = element.upper() - if normalized == "1A-C": - elements = ("1A", "1B", "1C") - else: - normalized = self._element(element) - elements = ("1sigma" if normalized == "1SIGMA" else normalized,) - fields = ("voltage", "current", "phase_angle", "real_power", "reactive_power") - return SnapshotStream(self, elements, fields) + 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, ...], fields: tuple[str, ...]) -> None: + def __init__(self, dev: OWH9830, elements: tuple[str, ...], max_harmonics: int = 0) -> None: self._dev = dev self._elements = elements - self._fields = fields + 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, Decimal | None]]: + async def __anext__(self) -> dict[str, dict[str, Any]]: if self._closed: raise StopAsyncIteration if not self._configured: @@ -288,11 +333,14 @@ async def __anext__(self) -> dict[str, dict[str, Decimal | None]]: raise RuntimeError("Snapshot iterator is already active") try: async with self._dev._io_lock: - if self._dev._snapshot_elements != self._elements: + 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) + await self._dev._configure_snapshot(self._elements, self._max_harmonics) self._dev._streaming_snapshots = True except BaseException: self._dev._streaming_snapshots = False @@ -302,17 +350,11 @@ async def __anext__(self) -> dict[str, dict[str, Decimal | None]]: async with self._dev._io_lock: try: - values = await self._dev._snapshot_values(self._elements) + return await self._dev._snapshot_values(self._elements, self._max_harmonics) except BaseException: self._dev._snapshot_elements = None await self.aclose() raise - return { - selected: dict( - zip(self._fields, values[index * len(self._fields) : (index + 1) * len(self._fields)], strict=True) - ) - for index, selected in enumerate(self._elements) - } async def aclose(self) -> None: if not self._closed: diff --git a/tests/test_owh9830.py b/tests/test_owh9830.py index f182bc7..a6b2836 100644 --- a/tests/test_owh9830.py +++ b/tests/test_owh9830.py @@ -329,3 +329,97 @@ async def test_measure_snapshots_stream() -> None: 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", + ] From 69bf757a5749c19a1fca8f3d6a4ad628af91252e Mon Sep 17 00:00:00 2001 From: Eero af Heurlin Date: Fri, 2 Oct 2026 15:17:06 +0300 Subject: [PATCH 12/12] docs(owh9830): document harmonic snapshot measurements and examples Co-authored-by: Junie --- examples/owh9830.rst | 66 +++++++++++++++++++++----------- examples/owh9830_async_stream.py | 21 ++++++++-- examples/owh9830_serial.py | 18 +++++++-- 3 files changed, 77 insertions(+), 28 deletions(-) diff --git a/examples/owh9830.rst b/examples/owh9830.rst index ff6d41e..04c2f78 100644 --- a/examples/owh9830.rst +++ b/examples/owh9830.rst @@ -50,35 +50,46 @@ from the same held set of meter-reported values:: 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). Single ``1A``, ``1B``, ``1C`` and ``1sigma`` selections -are supported, as is the combined ``1A-C`` selection. Inputs are case insensitive; -returned names are canonical. Values are ``Decimal`` or ``None`` for unavailable -readings (``----``), including phase/power on unconnected phases. +``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), then reuses it for repeated polls. It queries the -current HOLD state, enables HOLD if necessary, reads the configured slots, 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. +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 and update-rate commands sent through ``dev.command`` invalidate the -cached selection. +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 example is:: +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 ------------------- @@ -89,22 +100,33 @@ 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"]) -The iterator configures the numeric display page slots 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. + 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_async_stream.py /dev/tty.usbserial-A92TRYDJ --count 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 -------------------- diff --git a/examples/owh9830_async_stream.py b/examples/owh9830_async_stream.py index a103290..9b3fc8e 100644 --- a/examples/owh9830_async_stream.py +++ b/examples/owh9830_async_stream.py @@ -15,6 +15,13 @@ async def main() -> None: 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)" ) @@ -24,10 +31,18 @@ async def main() -> None: try: identity = await dev.identify() print(f"Connected to: {identity}") - print(f"Streaming snapshots for {args.element} (Ctrl-C to stop)...") - + 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 dev.measure_snapshots(args.element): + async for snapshot in stream: received += 1 print(f"[{received}]", snapshot) if args.count > 0 and received >= args.count: diff --git a/examples/owh9830_serial.py b/examples/owh9830_serial.py index 9b4f5aa..a6f11c6 100644 --- a/examples/owh9830_serial.py +++ b/examples/owh9830_serial.py @@ -14,6 +14,7 @@ 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)" ) @@ -24,14 +25,25 @@ if args.stream != 0: print(f"Streaming snapshots for {args.element} (Ctrl-C to stop)...") try: - for count, snapshot in enumerate(dev.measure_snapshots(args.element), 1): + 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: - print("Snapshot (V/A/degrees/W/var):", dev.measure_snapshot(args.element)) - if args.element in ("1A", "1B", "1C"): + 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"