#!/usr/bin/env python3 """ must_pv1800_monitor.py Real-time monitor for MUST PV1800 (2024) solar inverter/charger via Modbus RTU. Connects to a remote host (default 192.168.44.94, root, password "testing1") over SSH, opens an exclusive pyserial-based TCP proxy to /dev/ttyUSB0 on that host, polls the MUST PV1800 Modbus registers, and prints a live snapshot of CHARGER / INVERTER / SETTINGS values. Tested hardware: MUST PV1800 2024 variant (24V battery / 120V AC). Register map is shared across PH1800 / PV1800 / EP1800 / PV3500 / EP3500 per the manufacturer's Modbus RTU protocol doc v1.4.15. The proxy (also embedded below) is uploaded over SFTP and run as a detached background process; it pre-flushes any stale bytes sitting in the kernel TTY buffer and then bridges the UART to a localhost TCP port that we tunnel back through SSH. This is more reliable than socat because it sets TIOCEXCL on the TTY (which the PV1800 RS485 driver appears to require) and runs the stty-equivalent setup itself. Three transport modes are supported (auto-detected, first one that matches the CLI wins): --serial Open /dev/ttyUSB0 directly on the local host. --tcp Connect to a TCP bridge already running on the remote. --ssh (default) Upload the embedded proxy over SSH and tunnel. Usage: python3 must_pv1800_monitor.py # live updates, every 3 s python3 must_pv1800_monitor.py --once # one snapshot and exit python3 must_pv1800_monitor.py --interval 1 --json python3 must_pv1800_monitor.py --serial --serial-port /dev/ttyUSB0 References: - https://github.com/taHC81/MUST-ESPhome (esp8266-MUST-PV18.yaml) - PH1800/PV1800/EP1800/PV3500/EP3500 RS485 Modbus RTU communication Protocol 1.4.15 (register map embedded from the .xlsx in the repo) """ import argparse import os import select import socket import sys import time from dataclasses import dataclass from typing import Any, Callable, List, Optional, Tuple try: import serial # pyserial (only needed for direct serial mode) except ImportError: serial = None try: import paramiko # only needed for the auto-SSH transport except ImportError: paramiko = None # --------------------------------------------------------------------------- # Connection settings (defaults; can be overridden by CLI/env) # --------------------------------------------------------------------------- DEFAULT_SSH_HOST = "192.168.44.94" DEFAULT_SSH_USER = "root" DEFAULT_SSH_PASS = "testing1" DEFAULT_SSH_PORT = 22 DEFAULT_SERIAL_PORT = "/dev/ttyUSB0" DEFAULT_BAUDRATE = 19200 DEFAULT_SLAVE = 4 DEFAULT_POLL_INTERVAL = 3.0 # seconds # Inter-register delay. The PV1800 firmware does not always honour our # `count` parameter -- it sometimes dumps a 39- or 49-register block in # response to any read. If we ask for the next register too quickly the # leftover bytes from the previous dump pollute the next read. 250 ms # gives the inverter's UART ring buffer time to drain between requests. INTER_REGISTER_DELAY = 0.25 # How long to wait for a single Modbus RTU response. The PV1800 # sometimes dumps a 39-register block (83 bytes) when you asked for 1, # and at 19200 baud that's ~45 ms of air time. Allow 2 s so we get the # full oversized frame including any inter-register pauses. RESPONSE_TIMEOUT = 2.0 # --------------------------------------------------------------------------- # Modbus helpers # --------------------------------------------------------------------------- def _crc16(data: bytes) -> bytes: """Modbus RTU CRC16 (poly 0xA001, little-endian output).""" crc = 0xFFFF for b in data: crc ^= b for _ in range(8): crc = (crc >> 1) ^ 0xA001 if crc & 1 else crc >> 1 return bytes((crc & 0xFF, (crc >> 8) & 0xFF)) # ---- Transport abstraction ------------------------------------------------ class Transport: """Minimal interface: read/write/close + timeout + drain_input().""" def read_exact(self, n: int, timeout: float) -> bytes: raise NotImplementedError def write(self, data: bytes) -> int: raise NotImplementedError def drain_input(self) -> None: raise NotImplementedError def _stash_extra(self, data: bytes) -> None: """Buffer extra bytes that arrived after a complete response so the next read can pick them up instead of waiting for new data.""" # Default implementation: discard (most transports can't replay). # _SSHTunneledSerial overrides this. pass def close(self) -> None: raise NotImplementedError class _SerialTransport(Transport): def __init__(self, port: str, baudrate: int): if serial is None: raise RuntimeError("pyserial is not installed; " "pip install pyserial") self._ser = serial.Serial( port=port, baudrate=baudrate, bytesize=serial.EIGHTBITS, parity=serial.PARITY_NONE, stopbits=serial.STOPBITS_ONE, timeout=RESPONSE_TIMEOUT, write_timeout=RESPONSE_TIMEOUT, ) def drain_input(self) -> None: try: self._ser.reset_input_buffer() except Exception: pass def write(self, data: bytes) -> int: self._ser.write(data) self._ser.flush() return len(data) def read_exact(self, n: int, timeout: float) -> bytes: prev = self._ser.timeout self._ser.timeout = max(0.05, timeout) try: return self._ser.read(n) or b"" finally: self._ser.timeout = prev def close(self) -> None: try: self._ser.close() except Exception: pass class _SocketTransport(Transport): """TCP bridge to a serial port (e.g. socat/ser2net).""" def __init__(self, host: str, port: int): self._sock = socket.create_connection((host, port), timeout=10) self._sock.settimeout(RESPONSE_TIMEOUT) def drain_input(self) -> None: self._sock.setblocking(False) try: while True: data = self._sock.recv(4096) if not data: break except (BlockingIOError, socket.error): pass self._sock.setblocking(True) self._sock.settimeout(RESPONSE_TIMEOUT) def write(self, data: bytes) -> int: self._sock.sendall(data) return len(data) def read_exact(self, n: int, timeout: float) -> bytes: self._sock.settimeout(max(0.05, timeout)) chunks = bytearray() deadline = time.monotonic() + timeout while len(chunks) < n and time.monotonic() < deadline: try: chunk = self._sock.recv(n - len(chunks)) if not chunk: break chunks.extend(chunk) except socket.timeout: break except socket.error: break return bytes(chunks) def close(self) -> None: try: self._sock.shutdown(socket.SHUT_RDWR) except Exception: pass try: self._sock.close() except Exception: pass # --- Embedded remote TCP<->serial proxy ------------------------------------ # # We push this tiny script over SFTP and exec it on the remote host. It # opens the local /dev/ttyUSBX with pyserial and bridges it to a random # localhost TCP port. We then SSH-tunnel back to that port with a # direct-tcpip channel. # # Why a Python proxy and not socat? socat's rawer/raw modes don't # reliably drive the CH341 USB-serial adapter that ships with most # MUST PV1800 reference designs -- bytes get stuck in the kernel TTY # buffer and never reach the inverter. pyserial uses TIOCEXCL/UNEXCL # internally on open, which is what the inverter's RS485 driver is # actually waiting for. _REMOTE_PROXY_SCRIPT = r'''#!/usr/bin/env python3 """Tiny TCP<->serial bridge for MUST PV1800 monitoring. Listens on 127.0.0.1:$PROXY_PORT, forwards every byte to/from $PROXY_TTY using pyserial. Multiple sequential TCP clients are served one at a time. """ import os, sys, threading, time import serial, socket PORT = int(os.environ.get("PROXY_PORT", "9700")) TTY = os.environ.get("PROXY_TTY", "/dev/ttyUSB0") BAUD = int(os.environ.get("PROXY_BAUD", "19200")) def log(msg): print(f"[must-proxy] {msg}", file=sys.stderr, flush=True) try: # Pre-flush: read-and-discard whatever stale bytes are sitting in the # kernel TTY buffer from previous runs. Without this, those bytes get # delivered to the first client and corrupt its first read. pre = serial.Serial(TTY, BAUD, timeout=0.2) pre.reset_input_buffer() pre.reset_output_buffer() flushed = 0 quiet_for = 0.0 pre_start = time.monotonic() while time.monotonic() - pre_start < 3.0 and quiet_for < 0.3: chunk = pre.read(512) if chunk: flushed += len(chunk) quiet_for = 0.0 else: quiet_for += 0.05 time.sleep(0.05) pre.close() if flushed: log(f"flushed {flushed} stale bytes from kernel TTY buffer") ser = serial.Serial(TTY, BAUD, timeout=0.05) ser.reset_input_buffer() ser.reset_output_buffer() log(f"opened {TTY} @ {BAUD} fd={ser.fd}") except Exception as e: log(f"could not open {TTY}: {e}") sys.exit(1) def tty_to_tcp(client): while True: try: data = ser.read(256) if data: client.sendall(data) except Exception: return srv = socket.socket(socket.AF_INET, socket.SOCK_STREAM) srv.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) srv.bind(("127.0.0.1", PORT)) srv.listen(5) log(f"listening on 127.0.0.1:{PORT}") while True: client, addr = srv.accept() log(f"client {addr}") threading.Thread(target=tty_to_tcp, args=(client,), daemon=True).start() try: while True: data = client.recv(4096) if not data: break ser.write(data) ser.flush() except Exception as e: log(f"tcp err: {e}") try: client.close() except Exception: pass log(f"client closed") ''' class _SSHTunneledSerial(Transport): """ Auto-managed SSH transport: connects to a remote host, uploads a tiny pyserial-based TCP<->serial proxy via SFTP, executes it as a detached background process, and thereby exposes the remote /dev/ttyUSBX as a local TCP socket reachable through an SSH direct-tcpip channel. Cleaning up on close: the proxy process is killed and the uploaded proxy script is removed. """ PROXY_REMOTE_PATH = "/tmp/.must_pv1800_proxy.py" def __init__(self, ssh_host: str, ssh_user: str, ssh_pass: str, ssh_port: int, serial_port: str, baudrate: int): if paramiko is None: raise RuntimeError("paramiko is required for SSH transport; " "pip install paramiko") self._client = paramiko.SSHClient() self._client.set_missing_host_key_policy(paramiko.AutoAddPolicy()) self._client.connect( hostname=ssh_host, port=ssh_port, username=ssh_user, password=ssh_pass, allow_agent=False, look_for_keys=False, timeout=10, banner_timeout=10, auth_timeout=10, ) transport = self._client.get_transport() # Make sure no zombies hold /dev/ttyUSBX. Anything left over from # a previous crash will eat every Modbus response. ch_k = transport.open_session(timeout=10) ch_k.exec_command( f"pkill -9 -f 'socat.*{serial_port}' 2>/dev/null; " f"pkill -9 -f 'must_pv1800_proxy.py' 2>/dev/null; " f"pkill -9 -f 'cat {serial_port}' 2>/dev/null; " f"pkill -9 -f 'dd if={serial_port}' 2>/dev/null; " f"fuser -k {serial_port} 2>/dev/null; true" ) ch_k.recv_exit_status() time.sleep(0.3) # Make sure pyserial is installed on the remote host. ch_py = transport.open_session(timeout=30) ch_py.exec_command("python3 -c 'import serial' 2>/dev/null " "|| (apt-get install -y python3-serial 2>&1 " " | tail -3)") ch_py.recv_exit_status() try: ch_py.recv(8192) except Exception: pass # Push the proxy script over SFTP. sftp = self._client.open_sftp() try: with sftp.open(self.PROXY_REMOTE_PATH, "w") as f: f.write(_REMOTE_PROXY_SCRIPT) sftp.chmod(self.PROXY_REMOTE_PATH, 0o755) finally: sftp.close() # Pick a free TCP port on the remote side (9500-9600). self._remote_port = None for p in range(9500, 9600): ch_p = transport.open_session(timeout=5) ch_p.exec_command( f"ss -tlnH 'sport = :{p}' 2>/dev/null | grep -q LISTEN && " f"echo taken || echo free" ) ch_p.recv_exit_status() if b"free" in ch_p.recv(1024): self._remote_port = p break if self._remote_port is None: self._client.close() raise RuntimeError("could not find a free TCP port on remote host") # Launch the proxy detached from the SSH channel. We use a # non-looping command so the captured stdout is just the single # `ss` line we care about. ch_b = transport.open_session(timeout=15) ch_b.exec_command( f"PROXY_PORT={self._remote_port} PROXY_TTY={serial_port} " f"PROXY_BAUD={baudrate} setsid nohup python3 " f"{self.PROXY_REMOTE_PATH} /tmp/.must_pv1800_proxy.log 2>&1 & disown; " f"for i in 1 2 3 4 5 6 7 8 9 10; do " f" out=$(ss -tlnH 'sport = :{self._remote_port}'); " f" if [ -n \"$out\" ]; then echo \"$out\"; break; fi; " f" sleep 0.5; " f"done" ) ch_b.settimeout(10) listener_info = b"" try: while True: chunk = ch_b.recv(4096) if not chunk: break listener_info += chunk except Exception: pass try: ch_b.recv_exit_status() except Exception: pass listener_info = listener_info.decode(errors="replace") if "LISTEN" not in listener_info: self._client.close() raise RuntimeError( f"proxy did not bind 127.0.0.1:{self._remote_port} on " f"remote. listener_info={listener_info!r} " f"log: " + self._safe_remote_log(transport) ) # Open a direct-tcpip SSH channel from local -> remote 127.0.0.1:port. self._chan = transport.open_channel( "direct-tcpip", ("127.0.0.1", self._remote_port), ("127.0.0.1", 0), ) self._chan.settimeout(RESPONSE_TIMEOUT) self._stash = b"" @staticmethod def _safe_remote_log(transport) -> str: try: ch = transport.open_session(timeout=3) ch.exec_command("cat /tmp/.must_pv1800_proxy.log 2>&1 | tail -5") ch.recv_exit_status() return ch.recv(4096).decode(errors="replace") except Exception: return "(log unavailable)" def drain_input(self) -> None: """Drain everything pending on the channel, then wait for the channel to be QUIET for at least 100 ms before returning. This guarantees that no straggler bytes from a previous oversized inverter dump are still in flight when we issue the next request. The 100 ms quiet period was chosen to be larger than the time it takes the proxy to forward one full Modbus RTU frame at 19200 baud (a 100-byte frame at 8N1 is ~52 ms; 100 ms gives a safe margin). """ self._stash = b"" # First, non-blocking drain of anything immediately available. try: while True: r, _, _ = select.select([self._chan], [], [], 0) if not r: break self._chan.recv(4096) except Exception: pass # Then wait for the line to stay quiet for 100 ms. quiet_for = 0.0 deadline = time.monotonic() + 1.0 # max total drain time while quiet_for < 0.1 and time.monotonic() < deadline: r, _, _ = select.select([self._chan], [], [], 0.02) if r: try: self._chan.recv(4096) except Exception: break quiet_for = 0.0 else: quiet_for += 0.02 def _stash_extra(self, data: bytes) -> None: self._stash = getattr(self, "_stash", b"") + data def write(self, data: bytes) -> int: self._chan.sendall(data) return len(data) def read_exact(self, n: int, timeout: float) -> bytes: # Consume from the stash first (bytes the previous call didn't use). stash = getattr(self, "_stash", b"") if stash: take = stash[:n] self._stash = stash[n:] if len(take) >= n: return bytes(take) chunks = bytearray(take) else: chunks = bytearray() deadline = time.monotonic() + max(0.05, timeout) while len(chunks) < n and time.monotonic() < deadline: remaining = deadline - time.monotonic() r, _, _ = select.select([self._chan], [], [], remaining) if not r: break try: chunk = self._chan.recv(n - len(chunks)) if not chunk: break chunks.extend(chunk) except Exception: break return bytes(chunks) def close(self) -> None: try: self._chan.close() except Exception: pass try: transport = self._client.get_transport() ch = transport.open_session(timeout=5) ch.exec_command( f"pkill -9 -f 'must_pv1800_proxy.py' 2>/dev/null; " f"rm -f {self.PROXY_REMOTE_PATH}; true" ) ch.recv_exit_status() except Exception: pass try: self._client.close() except Exception: pass # ---- Modbus over Transport ----------------------------------------------- def modbus_read_holding(t: Transport, slave: int, address: int, count: int = 1, timeout: float = RESPONSE_TIMEOUT) -> List[int]: """ Send a Modbus function-3 (Read Holding Registers) request. Returns a list of unsigned 16-bit register values. Automatically retries on CRC mismatch or truncated responses -- the MUST PV1800 occasionally drops a byte at the start of a burst, so a single retry is enough to recover in practice. """ last_err = None for attempt in range(3): try: return _modbus_read_holding_once(t, slave, address, count, timeout) except IOError as e: last_err = e # Drain anything that came in for the failed attempt before # the next try so we don't feed it back into ourselves. try: t.drain_input() except Exception: pass time.sleep(0.05 * (attempt + 1)) raise last_err # type: ignore[misc] def _modbus_read_holding_once(t: Transport, slave: int, address: int, count: int, timeout: float) -> List[int]: """Single-shot Modbus read; raises IOError on failure. The PV1800 firmware doesn't honour our `count` field -- it often dumps a 39-register block when we asked for one. So instead of reading a fixed number of bytes we: 1. Read 5 bytes (the response header). 2. Parse the bytecount byte to find the actual frame length. 3. Read the rest of the frame. 4. Look for the response somewhere in the buffer (multiple frames may be concatenated if the inverter dumped more than one block's worth of data) by scanning for a CRC-valid frame. """ pdu = bytes((slave, 0x03, (address >> 8) & 0xFF, address & 0xFF, (count >> 8) & 0xFF, count & 0xFF)) request_frame = pdu + _crc16(pdu) t.drain_input() t.write(request_frame) expected_bc = 2 * count # Step 1: read at least the 5-byte header so we know the bytecount. buf = bytearray() deadline = time.monotonic() + timeout while len(buf) < 5 and time.monotonic() < deadline: chunk = t.read_exact(5 - len(buf), max(0.05, deadline - time.monotonic())) if chunk: buf.extend(chunk) if len(buf) < 5: raise IOError(f"short read: got {len(buf)} bytes, expected >= 5 " f"(buf={buf.hex()})") # Step 2: if the header looks like a valid response, read the rest. # The inverter's bc may legitimately be > 2*count (oversized dump). if (buf[0] == slave and buf[1] == 0x03 and 0 <= buf[2] <= 250): actual_bc = buf[2] target_len = 5 + actual_bc while len(buf) < target_len and time.monotonic() < deadline: chunk = t.read_exact(target_len - len(buf), max(0.05, deadline - time.monotonic())) if chunk: buf.extend(chunk) # Step 3: scan for a CRC-valid response inside the buffer. parsed = _find_valid_response(bytes(buf), slave, expected_bc) if parsed is None: raise IOError(f"no CRC-valid response in {len(buf)} bytes " f"(buf={buf.hex()})") offset, actual_bc, payload, valid_len = parsed # Step 4: stash any leftover bytes so the next call picks them up. leftover = bytes(buf[offset + valid_len:]) if leftover: t._stash_extra(leftover) return [int.from_bytes(payload[i:i + 2], "big") for i in range(0, expected_bc, 2)] def _find_valid_response(buf: bytes, slave: int, expected_bc: int) -> Optional[Tuple[int, int, bytes, int]]: """Scan `buf` for a CRC-valid Modbus RTU response from `slave`. Returns (offset, actual_bc, payload_bytes, total_len_consumed) on success, None if no valid response is found. Multiple frames may be concatenated in the buffer if the inverter sent them back-to-back. We try every byte alignment starting from offset 0 and pick the FIRST match with a valid CRC. The CRC check is what tells us we found a real frame boundary; raw header bytes alone are not enough because the inverter's large dumps can contain many coincidental `[04 03 bc]` substrings. """ for offset in range(len(buf) - 4): if buf[offset] != slave: continue func = buf[offset + 1] if func not in (0x03, 0x83): continue bc = buf[offset + 2] if bc < expected_bc or bc > 250: continue total = 5 + bc if offset + total > len(buf): continue frame = buf[offset:offset + total] if _crc16(frame[:3 + bc]) != frame[3 + bc:5 + bc]: continue payload = frame[3:3 + bc] return (offset, bc, payload, total) return None def u16(v: int) -> int: return v & 0xFFFF def s16(v: int) -> int: v &= 0xFFFF return v - 0x10000 if v & 0x8000 else v # --------------------------------------------------------------------------- # Decoded value tables (from PH1800/PV1800/EP1800/PV3500/EP3500 protocol v1.4.15) # --------------------------------------------------------------------------- CHARGER_WORKSTATE = { 0: "Initialization mode", 1: "Standby", 2: "Charging", 3: "Fault", 4: "Flash", } MPPT_STATE = { 0: "Stop", 1: "Start", 2: "Normal work", 3: "MPPT limit", 4: "Float", } CHARGING_STATE = { 0: "Stop", 1: "Precharge", 2: "CC", 3: "CV", 4: "Float", 5: "Current derating", } INVERTER_WORKSTATE = { 0: "PowerOn", 1: "SelfTest", 2: "OffGrid", 3: "GridTie", 4: "ByPass", 5: "Stop", 6: "GridCharging", } BATTERY_TYPE = { 0: "AGM", 1: "Flooded", 2: "User defined battery", 3: "Lithium", } ENERGY_USE_MODE = { 0: "SBU (Solar/battery/utility)", 1: "SBU (Solar/battery/utility)", 2: "SUB (Solar/utility/battery)", 3: "UTI (Utility only)", 4: "SOL (Solar only)", } SOLARUSE_AIM = { 0: "LBU (Less battery use)", 1: "LBU (Less battery use)", 2: "OSO (Only solar output)", } # --------------------------------------------------------------------------- # Register map (from MUST-ESPhome ESPHome config) # --------------------------------------------------------------------------- @dataclass(frozen=True) class Register: address: int name: str scale: float = 1.0 signed: bool = False decode: Optional[Callable[[int], Any]] = None unit: str = "" # CHARGER block CHARGER_REGS: List[Register] = [ Register(15201, "Charger workstate", decode=lambda v: CHARGER_WORKSTATE.get(v, f"unknown({v})")), Register(15202, "MPPT state", decode=lambda v: MPPT_STATE.get(v, f"unknown({v})")), Register(15203, "Charging state", decode=lambda v: CHARGING_STATE.get(v, f"unknown({v})")), Register(15205, "PV voltage", scale=0.1, unit="V"), Register(15206, "Battery voltage", scale=0.1, unit="V"), Register(15207, "Charger current", scale=0.1, unit="A"), Register(15208, "Charger power", unit="W"), # 15211 Battery Relay, 15212 PV Relay -- not exposed by reference config ] # INVERTER block (read-only telemetry) INVERTER_REGS: List[Register] = [ Register(25201, "Inverter Work state", decode=lambda v: INVERTER_WORKSTATE.get(v, f"unknown({v})")), Register(25205, "Battery voltage", scale=0.1, unit="V"), Register(25206, "Inverter voltage", scale=0.1, unit="V"), Register(25207, "Grid voltage", scale=0.1, unit="V"), Register(25213, "Inverter power", signed=True, unit="W"), Register(25214, "Grid power", signed=True, unit="W"), Register(25215, "Load power", signed=True, unit="W"), Register(25216, "System load", unit="%"), Register(25233, "AC radiator temp", unit="°C"), Register(25234, "Transformer temp", unit="°C"), Register(25235, "DC radiator temp", unit="°C"), Register(25273, "Battery power", signed=True, unit="W"), Register(25274, "Battery current", signed=True, scale=0.1, unit="A"), ] # SETTINGS block (read-only mirror of writable config) SETTING_REGS: List[Register] = [ Register(10103, "Float voltage", scale=0.1, unit="V"), Register(10104, "Absorb voltage", scale=0.1, unit="V"), Register(20118, "Battery stop discharging voltage", scale=0.1, unit="V"), Register(20119, "Battery stop charging voltage", scale=0.1, unit="V"), Register(20127, "Battery low voltage", scale=0.1, unit="V"), Register(20128, "Battery high voltage", scale=0.1, unit="V"), Register(20132, "Charger current", scale=0.1, unit="A"), ] # --------------------------------------------------------------------------- # Polling / decoding # --------------------------------------------------------------------------- def poll_block(t: Transport, slave: int, regs: List[Register], base_addr: int) -> List[Tuple[str, Any, str]]: """ Read the entire contiguous register block starting at `base_addr` in ONE Modbus read, then pluck out the registers we actually care about from the response. The PV1800 firmware ignores our `count` anyway and dumps the whole block, so doing this explicitly makes the request fast AND unambiguous -- we can identify which frame in a multi-frame buffer is ours because the response contains the entire block starting at base_addr. `regs` should be a list of `Register` whose addresses all fall in the [base_addr, base_addr + count - 1] range. The function returns (name, decoded_value, unit) tuples for every register in `regs`, in the same order. """ if not regs: return [] block_len = max(r.address for r in regs) - base_addr + 1 raw_regs = modbus_read_holding(t, slave, base_addr, count=block_len) # raw_regs[0] == base_addr, raw_regs[i] == base_addr + i. results: List[Tuple[str, Any, str]] = [] for reg in regs: idx = reg.address - base_addr raw = raw_regs[idx] if reg.decode is not None: value = reg.decode(raw) elif reg.signed: value = s16(raw) else: value = raw * reg.scale results.append((reg.name, value, reg.unit)) return results def poll_individual(t: Transport, slave: int, regs: List[Register]) -> List[Tuple[str, Any, str]]: """ Read each register individually with a small inter-register delay. Used for settings which are scattered across non-contiguous ranges. """ results: List[Tuple[str, Any, str]] = [] for reg in regs: raw = modbus_read_holding(t, slave, reg.address, count=1)[0] if reg.decode is not None: value = reg.decode(raw) elif reg.signed: value = s16(raw) else: value = raw * reg.scale results.append((reg.name, value, reg.unit)) time.sleep(INTER_REGISTER_DELAY) return results def print_snapshot(snapshot: List[Tuple[str, Any, str]], title: str) -> None: bar = "-" * max(15, len(title) + 4) print(bar) print(f" {title}") print(bar) name_w = max(len(n) for n, _, _ in snapshot) for name, value, unit in snapshot: if isinstance(value, float): vstr = f"{value:.1f}" else: vstr = str(value) suffix = f"{unit}" if unit else "" print(f" {name:<{name_w}} = {vstr}{suffix}") print() # CLI # --------------------------------------------------------------------------- def parse_args() -> argparse.Namespace: p = argparse.ArgumentParser( description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter, ) g = p.add_mutually_exclusive_group() g.add_argument("--serial", dest="mode", action="store_const", const="serial", help="Open --serial-port locally (default if --tcp/--ssh absent).") g.add_argument("--tcp", dest="mode", action="store_const", const="tcp", help="Connect to a TCP bridge (socat/ser2net).") g.add_argument("--ssh", dest="mode", action="store_const", const="ssh", help="Tunnel over SSH and start a remote pyserial-based proxy (default).") p.set_defaults(mode="ssh") p.add_argument("--tcp-host", default=os.environ.get("MUST_TCP_HOST", "127.0.0.1")) p.add_argument("--tcp-port", type=int, default=int(os.environ.get("MUST_TCP_PORT", "8502"))) p.add_argument("--ssh-host", default=os.environ.get("MUST_SSH_HOST", DEFAULT_SSH_HOST)) p.add_argument("--ssh-port", type=int, default=int(os.environ.get("MUST_SSH_PORT", DEFAULT_SSH_PORT))) p.add_argument("--ssh-user", default=os.environ.get("MUST_SSH_USER", DEFAULT_SSH_USER)) p.add_argument("--ssh-pass", default=os.environ.get("MUST_SSH_PASS", DEFAULT_SSH_PASS)) p.add_argument("--serial-port", default=os.environ.get("MUST_SERIAL_PORT", DEFAULT_SERIAL_PORT)) p.add_argument("--baudrate", type=int, default=int(os.environ.get("MUST_BAUD", DEFAULT_BAUDRATE))) p.add_argument("--slave", type=int, default=int(os.environ.get("MUST_SLAVE", DEFAULT_SLAVE))) p.add_argument("--interval", type=float, default=float(os.environ.get("MUST_INTERVAL", DEFAULT_POLL_INTERVAL)), help="Polling interval in seconds (default: %(default)s)") p.add_argument("--once", action="store_true", help="Run a single snapshot and exit") p.add_argument("--json", action="store_true", help="Print each snapshot as a JSON object (one line per snapshot).") return p.parse_args() def open_transport(args) -> Tuple[Transport, str]: """ Open the requested transport and return (transport, label). """ if args.mode == "serial": return _SerialTransport(args.serial_port, args.baudrate), \ f"serial {args.serial_port} @ {args.baudrate}" if args.mode == "tcp": return _SocketTransport(args.tcp_host, args.tcp_port), \ f"tcp {args.tcp_host}:{args.tcp_port}" # default: ssh return _SSHTunneledSerial( ssh_host=args.ssh_host, ssh_user=args.ssh_user, ssh_pass=args.ssh_pass, ssh_port=args.ssh_port, serial_port=args.serial_port, baudrate=args.baudrate, ), f"ssh {args.ssh_user}@{args.ssh_host} -> {args.serial_port}" def snapshot_to_dicts(charger, inverter, settings) -> dict: """Combine three (name,value,unit) lists into a single dict for JSON.""" out = {} for src in (charger, inverter, settings): for name, value, unit in src: if isinstance(value, float): out[name] = {"value": round(value, 3), "unit": unit} else: out[name] = {"value": value, "unit": unit} return out def main() -> int: args = parse_args() print(f"Connecting to MUST PV1800 ({args.mode}): " f"slave={args.slave}, interval={args.interval}s", file=sys.stderr) try: t, label = open_transport(args) except Exception as e: print(f"ERROR: could not open transport: {e}", file=sys.stderr) return 1 print(f"Transport: {label}", file=sys.stderr) try: # Quick connectivity check: read Charger workstate register. try: sanity = modbus_read_holding(t, args.slave, 15201, count=1, timeout=2.0) print(f"Sanity OK: register 15201 = {sanity[0]} " f"({CHARGER_WORKSTATE.get(sanity[0], '?')})", file=sys.stderr) except Exception as e: print(f"ERROR: initial Modbus read failed: {e}", file=sys.stderr) print("Hints:", file=sys.stderr) print(" - Confirm /dev/ttyUSB0 exists on the remote host.", file=sys.stderr) print(" - Confirm baud rate (default 19200, 8N1).", file=sys.stderr) print(" - Confirm slave ID (default 4).", file=sys.stderr) print(" - Another process might already hold the port " "(check `fuser /dev/ttyUSB0` over SSH).", file=sys.stderr) return 2 next_tick = time.monotonic() attempts = 0 max_once_attempts = 5 # for --once: try this many times before giving up while True: now = time.strftime("%Y-%m-%d %H:%M:%S") try: # CHARGER block: reads registers 15201..15208 (8 regs). # We always start at 15201 even if the firmware returns # more -- the inverter ignores our count and dumps the # whole block, so this is the unambiguous read. charger = poll_block(t, args.slave, CHARGER_REGS, base_addr=15201) # INVERTER block: reads registers 25201..25274 (74 regs). inverter = poll_block(t, args.slave, INVERTER_REGS, base_addr=25201) # SETTINGS are split across two non-contiguous ranges: # 10103..10104 (2 regs) and 20118..20132 (15 regs). # Polling them as blocks avoids the per-register cross- # frame mix-up we get from individual reads. s1 = poll_block(t, args.slave, [r for r in SETTING_REGS if r.address < 11000], base_addr=10103) s2 = poll_block(t, args.slave, [r for r in SETTING_REGS if r.address >= 11000], base_addr=20118) settings = s1 + s2 except Exception as e: attempts += 1 print(f"[{now}] poll error: {e}", file=sys.stderr) if args.once and attempts >= max_once_attempts: return 3 time.sleep(args.interval) continue attempts = 0 if args.json: import json payload = {"timestamp": now, "snapshot": snapshot_to_dicts( charger, inverter, settings)} print(json.dumps(payload, ensure_ascii=False)) else: print(f"=== {now} ===") print_snapshot(charger, "CHARGER") print_snapshot(inverter, "INVERTER") print_snapshot(settings, "SETTINGS") sys.stdout.flush() if args.once: return 0 next_tick += args.interval sleep_for = next_tick - time.monotonic() if sleep_for > 0: time.sleep(sleep_for) else: next_tick = time.monotonic() finally: try: t.close() except Exception: pass if __name__ == "__main__": sys.exit(main())