| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024 |
- #!/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())
|