|
@@ -0,0 +1,1024 @@
|
|
|
|
|
+#!/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} </dev/null "
|
|
|
|
|
+ f">/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())
|