| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197 |
- #!/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
- # ---------------------------------------------------------------------------
- # Cross-checked against the official protocol document
- # `PH1800 PV1800 EP1800 PV3500 EP3500 RS485 Modbus RTU communication
- # Protocol 1.4.15.xlsx` (taHC81/MUST-ESPhome repo; the same document
- # shared on Scribd as document 881604771).
- CHARGER_WORKSTATE = {
- 0: "Initialization mode",
- 1: "SelfTest mode",
- 2: "Work mode",
- 3: "Stop mode",
- }
- MPPT_STATE = {
- 0: "Stop",
- 1: "MPPT",
- 2: "Current limiting",
- }
- CHARGING_STATE = {
- 0: "Stop",
- 1: "Absorb charge",
- 2: "Float charge",
- 3: "EQ charge",
- }
- INVERTER_WORKSTATE = {
- 0: "PowerOn",
- 1: "SelfTest",
- 2: "OffGrid",
- 3: "GridTie",
- 4: "ByPass",
- 5: "Stop",
- 6: "GridCharging",
- }
- def RELAY_STATE(v: int) -> str:
- """Decode relay-state registers (25237..25242, 15211, 15212).
- Per the protocol: 0=Disconnect, 1=Connect.
- """
- return "Closed" if v == 1 else "Open"
- 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 (registers 15201..15221; gaps at 15204, 15213, 15214
- # are reserved / unused by the firmware)
- CHARGER_REGS: List[Register] = [
- # ----- Operating state -----
- 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})")),
- # ----- Live measurements -----
- Register(15205, "PV voltage", scale=0.1, unit="V"),
- Register(15206, "Battery voltage (charger)", scale=0.1, unit="V"),
- Register(15207, "Charger current", scale=0.1, unit="A"),
- Register(15208, "Charger power", unit="W"),
- # ----- Sensors -----
- Register(15209, "Charger radiator temp", unit="°C"),
- Register(15210, "External temp", unit="°C"),
- # ----- Relays -----
- Register(15211, "Battery Relay (charger)", decode=RELAY_STATE),
- Register(15212, "PV Relay (charger)", decode=RELAY_STATE),
- # ----- Configuration / identity -----
- Register(15215, "Battery voltage grade", decode=lambda v: f"{v}V system"),
- Register(15216, "Charger rated current", scale=0.1, unit="A"),
- # ----- Energy accumulator (high * 1000 kWh + low * 0.1 kWh) -----
- Register(15217, "Acc PV energy MWh part", unit="MWh part"),
- Register(15218, "Acc PV energy kWh part", scale=0.1, unit="kWh part"),
- # ----- Runtime since last reset -----
- Register(15219, "Acc runtime days", unit="d"),
- Register(15220, "Acc runtime hours", unit="h"),
- Register(15221, "Acc runtime minutes", unit="m"),
- ]
- # INVERTER block (registers 25201..25274 with gaps at 25204, 25220,
- # 25224, 25227, 25228, 25232, 25236, 25243, 25244, 25262..25268, 25270..25272,
- # 25275, 25276).
- INVERTER_REGS: List[Register] = [
- # ----- Operating state -----
- Register(25201, "Inverter Work state", decode=lambda v: INVERTER_WORKSTATE.get(v, f"unknown({v})")),
- Register(25202, "AC voltage grade", decode=lambda v: f"{v}V variant"),
- Register(25203, "Rated power", unit="VA"),
- # ----- Voltages -----
- Register(25205, "Battery voltage (inverter)", scale=0.1, unit="V"),
- Register(25206, "Inverter voltage", scale=0.1, unit="V"),
- Register(25207, "Grid voltage", scale=0.1, unit="V"),
- Register(25208, "BUS voltage", scale=0.1, unit="V"),
- # ----- Currents -----
- Register(25209, "Control current", scale=0.1, unit="A"),
- Register(25210, "Inverter current", scale=0.1, unit="A"),
- Register(25211, "Grid current", scale=0.1, unit="A"),
- Register(25212, "Load current", scale=0.1, unit="A"),
- # ----- Real power (P) -----
- Register(25213, "PInverter", signed=True, unit="W"),
- Register(25214, "PGrid", signed=True, unit="W"),
- Register(25215, "PLoad", signed=True, unit="W"),
- Register(25216, "Load percent", unit="%"),
- # ----- Apparent power (S) -----
- Register(25217, "SInverter", signed=True, unit="VA"),
- Register(25218, "SGrid", signed=True, unit="VA"),
- Register(25219, "SLoad", signed=True, unit="VA"),
- # ----- Reactive power (Q) -----
- Register(25221, "QInverter", signed=True, unit="var"),
- Register(25222, "QGrid", signed=True, unit="var"),
- Register(25223, "QLoad", signed=True, unit="var"),
- # ----- Frequencies -----
- Register(25225, "Inverter frequency", scale=0.01, unit="Hz"),
- Register(25226, "Grid frequency", scale=0.01, unit="Hz"),
- # ----- Temperatures -----
- Register(25233, "AC radiator temp", unit="°C"),
- Register(25234, "Transformer temp", unit="°C"),
- Register(25235, "DC radiator temp", unit="°C"),
- # ----- Relays -----
- Register(25237, "Inverter relay", decode=RELAY_STATE),
- Register(25238, "Grid relay", decode=RELAY_STATE),
- Register(25239, "Load relay", decode=RELAY_STATE),
- Register(25240, "N_Line relay", decode=RELAY_STATE),
- Register(25241, "DC relay", decode=RELAY_STATE),
- Register(25242, "Earth relay", decode=RELAY_STATE),
- # ----- Energy accumulators (each = high*1000 kWh + low*0.1 kWh) -----
- Register(25245, "Acc charger MWh part", unit="MWh part"),
- Register(25246, "Acc charger kWh part", scale=0.1, unit="kWh part"),
- Register(25247, "Acc discharger MWh part", unit="MWh part"),
- Register(25248, "Acc discharger kWh part", scale=0.1, unit="kWh part"),
- Register(25249, "Acc grid-buy MWh part", unit="MWh part"),
- Register(25250, "Acc grid-buy kWh part", scale=0.1, unit="kWh part"),
- Register(25251, "Acc grid-sell MWh part", unit="MWh part"),
- Register(25252, "Acc grid-sell kWh part", scale=0.1, unit="kWh part"),
- Register(25253, "Acc load MWh part", unit="MWh part"),
- Register(25254, "Acc load kWh part", scale=0.1, unit="kWh part"),
- Register(25255, "Acc self-use MWh part", unit="MWh part"),
- Register(25256, "Acc self-use kWh part", scale=0.1, unit="kWh part"),
- Register(25257, "Acc PV-sell MWh part", unit="MWh part"),
- Register(25258, "Acc PV-sell kWh part", scale=0.1, unit="kWh part"),
- Register(25259, "Acc grid-charge MWh part", unit="MWh part"),
- Register(25260, "Acc grid-charge kWh part", scale=0.1, unit="kWh part"),
- # ----- Battery telemetry (live) -----
- Register(25273, "Battery power", signed=True, unit="W"),
- Register(25274, "Battery current", signed=True, scale=0.1, unit="A"),
- Register(25275, "Battery voltage grade", decode=lambda v: f"{v}V system"),
- ]
- # 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 _split_groups(snapshot: List[Tuple[str, Any, str]],
- headings: List[Tuple[str, Any]]) -> List[Tuple[str, List[Tuple[str, Any, str]]]]:
- """Split a flat snapshot into labelled groups.
- `headings` is a list of `(heading_text, name_prefix_or_prefixes)`
- pairs. Each register whose name starts with one of the listed
- prefixes goes into that group. A prefix can also be a `str` for
- the common single-prefix case. Registers that don't match any
- prefix are silently dropped.
- """
- out: List[Tuple[str, List[Tuple[str, Any, str]]]] = []
- for header, prefixes in headings:
- if isinstance(prefixes, str):
- prefixes = [prefixes]
- items = [(n, v, u) for n, v, u in snapshot
- if any(n.startswith(p) for p in prefixes)]
- if items:
- out.append((header, items))
- return out
- def _combine_accumulators(snapshot: List[Tuple[str, Any, str]]) -> List[Tuple[str, Any, str]]:
- """Merge every 'Acc <name> MWh part' / 'kWh part' pair into a single
- 'Acc <name> (total) = X.X kWh' line. Returns the merged list (in
- the original order, with the combined line replacing the two
- originals)."""
- by_name = {n: (v, u) for n, v, u in snapshot}
- out: List[Tuple[str, Any, str]] = []
- skip = set()
- for name, value, unit in snapshot:
- if name in skip:
- continue
- if name.endswith(" MWh part"):
- base = name[: -len(" MWh part")]
- kwh_name = base + " kWh part"
- if kwh_name in by_name:
- kwh_val = by_name[kwh_name][0]
- total_kwh = value * 1000.0 + kwh_val * 0.1
- # `base` looks like "Acc <name>"; we want "<name>".
- if base.startswith("Acc "):
- label = base[len("Acc "):]
- else:
- label = base
- out.append((f"Acc {label.strip()} (total)",
- round(total_kwh, 2), "kWh"))
- skip.add(name)
- skip.add(kwh_name)
- continue
- out.append((name, value, unit))
- return out
- def print_snapshot(snapshot: List[Tuple[str, Any, str]], title: str,
- groups: Optional[List[Tuple[str, List[Tuple[str, Any, str]]]]] = None) -> None:
- """Print a snapshot. If `groups` is given, print section headers between
- groups; otherwise dump the whole list flat."""
- bar = "-" * max(15, len(title) + 4)
- print(bar)
- print(f" {title}")
- print(bar)
- name_w = max(len(n) for n, _, _ in snapshot)
- if groups:
- for header, items in groups:
- print(f" --- {header} ---")
- for name, value, unit in items:
- 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()
- else:
- 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} ===")
- # Combine accumulator pairs (MWh part + kWh part -> total kWh)
- # and split CHARGER/INVERTER into readable sub-sections.
- charger = _combine_accumulators(charger)
- inverter = _combine_accumulators(inverter)
- charger_groups = _split_groups(charger, [
- ("State", ["Charger workstate", "MPPT state", "Charging state"]),
- ("Live readings", ["PV voltage", "Battery voltage (charger)", "Charger current", "Charger power"]),
- ("Sensors", ["Charger radiator temp", "External temp"]),
- ("Relays", ["Battery Relay (charger)", "PV Relay (charger)"]),
- ("Configuration", ["Battery voltage grade", "Charger rated current"]),
- ("Energy totals", ["Acc PV energy"]),
- ("Runtime", ["Acc runtime "]),
- ])
- inverter_groups = _split_groups(inverter, [
- ("State", ["Inverter Work state"]),
- ("Configuration", ["AC voltage grade", "Rated power", "Battery voltage grade"]),
- ("Voltages", ["Battery voltage (inverter)", "Inverter voltage", "Grid voltage", "BUS voltage"]),
- ("Currents", ["Control current", "Inverter current", "Grid current", "Load current"]),
- ("Real power (P)", ["PInverter", "PGrid", "PLoad", "Load percent"]),
- ("Apparent power (S)", ["SInverter", "SGrid", "SLoad"]),
- ("Reactive power (Q)", ["QInverter", "QGrid", "QLoad"]),
- ("Frequencies", ["Inverter frequency", "Grid frequency"]),
- ("Temperatures", ["AC radiator temp", "Transformer temp", "DC radiator temp"]),
- ("Relays", ["Inverter relay", "Grid relay", "Load relay", "N_Line relay", "DC relay", "Earth relay"]),
- ("Energy totals", ["Acc "]),
- ("Battery telemetry", ["Battery power", "Battery current"]),
- ])
- print_snapshot(charger, "CHARGER", groups=charger_groups)
- print_snapshot(inverter, "INVERTER", groups=inverter_groups)
- 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())
|