must_pv1800_monitor.py 37 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023
  1. #!/usr/bin/env python3
  2. """
  3. must_pv1800_monitor.py
  4. Real-time monitor for MUST PV1800 (2024) solar inverter/charger via Modbus RTU.
  5. Connects to a remote host (default 192.168.44.94, root, password
  6. "testing1") over SSH, opens an exclusive pyserial-based TCP proxy to
  7. /dev/ttyUSB0 on that host, polls the MUST PV1800 Modbus registers, and
  8. prints a live snapshot of CHARGER / INVERTER / SETTINGS values.
  9. Tested hardware: MUST PV1800 2024 variant (24V battery / 120V AC).
  10. Register map is shared across PH1800 / PV1800 / EP1800 / PV3500 / EP3500
  11. per the manufacturer's Modbus RTU protocol doc v1.4.15.
  12. The proxy (also embedded below) is uploaded over SFTP and run as a
  13. detached background process; it pre-flushes any stale bytes sitting in
  14. the kernel TTY buffer and then bridges the UART to a localhost TCP port
  15. that we tunnel back through SSH. This is more reliable than socat
  16. because it sets TIOCEXCL on the TTY (which the PV1800 RS485 driver
  17. appears to require) and runs the stty-equivalent setup itself.
  18. Three transport modes are supported (auto-detected, first one that
  19. matches the CLI wins):
  20. --serial Open /dev/ttyUSB0 directly on the local host.
  21. --tcp Connect to a TCP bridge already running on the remote.
  22. --ssh (default) Upload the embedded proxy over SSH and tunnel.
  23. Usage:
  24. python3 must_pv1800_monitor.py # live updates, every 3 s
  25. python3 must_pv1800_monitor.py --once # one snapshot and exit
  26. python3 must_pv1800_monitor.py --interval 1 --json
  27. python3 must_pv1800_monitor.py --serial --serial-port /dev/ttyUSB0
  28. References:
  29. - https://github.com/taHC81/MUST-ESPhome (esp8266-MUST-PV18.yaml)
  30. - PH1800/PV1800/EP1800/PV3500/EP3500 RS485 Modbus RTU communication
  31. Protocol 1.4.15 (register map embedded from the .xlsx in the repo)
  32. """
  33. import argparse
  34. import os
  35. import select
  36. import socket
  37. import sys
  38. import time
  39. from dataclasses import dataclass
  40. from typing import Any, Callable, List, Optional, Tuple
  41. try:
  42. import serial # pyserial (only needed for direct serial mode)
  43. except ImportError:
  44. serial = None
  45. try:
  46. import paramiko # only needed for the auto-SSH transport
  47. except ImportError:
  48. paramiko = None
  49. # ---------------------------------------------------------------------------
  50. # Connection settings (defaults; can be overridden by CLI/env)
  51. # ---------------------------------------------------------------------------
  52. DEFAULT_SSH_HOST = "192.168.44.94"
  53. DEFAULT_SSH_USER = "root"
  54. DEFAULT_SSH_PASS = "testing1"
  55. DEFAULT_SSH_PORT = 22
  56. DEFAULT_SERIAL_PORT = "/dev/ttyUSB0"
  57. DEFAULT_BAUDRATE = 19200
  58. DEFAULT_SLAVE = 4
  59. DEFAULT_POLL_INTERVAL = 3.0 # seconds
  60. # Inter-register delay. The PV1800 firmware does not always honour our
  61. # `count` parameter -- it sometimes dumps a 39- or 49-register block in
  62. # response to any read. If we ask for the next register too quickly the
  63. # leftover bytes from the previous dump pollute the next read. 250 ms
  64. # gives the inverter's UART ring buffer time to drain between requests.
  65. INTER_REGISTER_DELAY = 0.25
  66. # How long to wait for a single Modbus RTU response. The PV1800
  67. # sometimes dumps a 39-register block (83 bytes) when you asked for 1,
  68. # and at 19200 baud that's ~45 ms of air time. Allow 2 s so we get the
  69. # full oversized frame including any inter-register pauses.
  70. RESPONSE_TIMEOUT = 2.0
  71. # ---------------------------------------------------------------------------
  72. # Modbus helpers
  73. # ---------------------------------------------------------------------------
  74. def _crc16(data: bytes) -> bytes:
  75. """Modbus RTU CRC16 (poly 0xA001, little-endian output)."""
  76. crc = 0xFFFF
  77. for b in data:
  78. crc ^= b
  79. for _ in range(8):
  80. crc = (crc >> 1) ^ 0xA001 if crc & 1 else crc >> 1
  81. return bytes((crc & 0xFF, (crc >> 8) & 0xFF))
  82. # ---- Transport abstraction ------------------------------------------------
  83. class Transport:
  84. """Minimal interface: read/write/close + timeout + drain_input()."""
  85. def read_exact(self, n: int, timeout: float) -> bytes: raise NotImplementedError
  86. def write(self, data: bytes) -> int: raise NotImplementedError
  87. def drain_input(self) -> None: raise NotImplementedError
  88. def _stash_extra(self, data: bytes) -> None:
  89. """Buffer extra bytes that arrived after a complete response so the
  90. next read can pick them up instead of waiting for new data."""
  91. # Default implementation: discard (most transports can't replay).
  92. # _SSHTunneledSerial overrides this.
  93. pass
  94. def close(self) -> None: raise NotImplementedError
  95. class _SerialTransport(Transport):
  96. def __init__(self, port: str, baudrate: int):
  97. if serial is None:
  98. raise RuntimeError("pyserial is not installed; "
  99. "pip install pyserial")
  100. self._ser = serial.Serial(
  101. port=port, baudrate=baudrate,
  102. bytesize=serial.EIGHTBITS, parity=serial.PARITY_NONE,
  103. stopbits=serial.STOPBITS_ONE,
  104. timeout=RESPONSE_TIMEOUT, write_timeout=RESPONSE_TIMEOUT,
  105. )
  106. def drain_input(self) -> None:
  107. try:
  108. self._ser.reset_input_buffer()
  109. except Exception:
  110. pass
  111. def write(self, data: bytes) -> int:
  112. self._ser.write(data)
  113. self._ser.flush()
  114. return len(data)
  115. def read_exact(self, n: int, timeout: float) -> bytes:
  116. prev = self._ser.timeout
  117. self._ser.timeout = max(0.05, timeout)
  118. try:
  119. return self._ser.read(n) or b""
  120. finally:
  121. self._ser.timeout = prev
  122. def close(self) -> None:
  123. try:
  124. self._ser.close()
  125. except Exception:
  126. pass
  127. class _SocketTransport(Transport):
  128. """TCP bridge to a serial port (e.g. socat/ser2net)."""
  129. def __init__(self, host: str, port: int):
  130. self._sock = socket.create_connection((host, port), timeout=10)
  131. self._sock.settimeout(RESPONSE_TIMEOUT)
  132. def drain_input(self) -> None:
  133. self._sock.setblocking(False)
  134. try:
  135. while True:
  136. data = self._sock.recv(4096)
  137. if not data:
  138. break
  139. except (BlockingIOError, socket.error):
  140. pass
  141. self._sock.setblocking(True)
  142. self._sock.settimeout(RESPONSE_TIMEOUT)
  143. def write(self, data: bytes) -> int:
  144. self._sock.sendall(data)
  145. return len(data)
  146. def read_exact(self, n: int, timeout: float) -> bytes:
  147. self._sock.settimeout(max(0.05, timeout))
  148. chunks = bytearray()
  149. deadline = time.monotonic() + timeout
  150. while len(chunks) < n and time.monotonic() < deadline:
  151. try:
  152. chunk = self._sock.recv(n - len(chunks))
  153. if not chunk:
  154. break
  155. chunks.extend(chunk)
  156. except socket.timeout:
  157. break
  158. except socket.error:
  159. break
  160. return bytes(chunks)
  161. def close(self) -> None:
  162. try:
  163. self._sock.shutdown(socket.SHUT_RDWR)
  164. except Exception:
  165. pass
  166. try:
  167. self._sock.close()
  168. except Exception:
  169. pass
  170. # --- Embedded remote TCP<->serial proxy ------------------------------------
  171. #
  172. # We push this tiny script over SFTP and exec it on the remote host. It
  173. # opens the local /dev/ttyUSBX with pyserial and bridges it to a random
  174. # localhost TCP port. We then SSH-tunnel back to that port with a
  175. # direct-tcpip channel.
  176. #
  177. # Why a Python proxy and not socat? socat's rawer/raw modes don't
  178. # reliably drive the CH341 USB-serial adapter that ships with most
  179. # MUST PV1800 reference designs -- bytes get stuck in the kernel TTY
  180. # buffer and never reach the inverter. pyserial uses TIOCEXCL/UNEXCL
  181. # internally on open, which is what the inverter's RS485 driver is
  182. # actually waiting for.
  183. _REMOTE_PROXY_SCRIPT = r'''#!/usr/bin/env python3
  184. """Tiny TCP<->serial bridge for MUST PV1800 monitoring.
  185. Listens on 127.0.0.1:$PROXY_PORT, forwards every byte to/from $PROXY_TTY
  186. using pyserial. Multiple sequential TCP clients are served one at a time.
  187. """
  188. import os, sys, threading, time
  189. import serial, socket
  190. PORT = int(os.environ.get("PROXY_PORT", "9700"))
  191. TTY = os.environ.get("PROXY_TTY", "/dev/ttyUSB0")
  192. BAUD = int(os.environ.get("PROXY_BAUD", "19200"))
  193. def log(msg):
  194. print(f"[must-proxy] {msg}", file=sys.stderr, flush=True)
  195. try:
  196. # Pre-flush: read-and-discard whatever stale bytes are sitting in the
  197. # kernel TTY buffer from previous runs. Without this, those bytes get
  198. # delivered to the first client and corrupt its first read.
  199. pre = serial.Serial(TTY, BAUD, timeout=0.2)
  200. pre.reset_input_buffer()
  201. pre.reset_output_buffer()
  202. flushed = 0
  203. quiet_for = 0.0
  204. pre_start = time.monotonic()
  205. while time.monotonic() - pre_start < 3.0 and quiet_for < 0.3:
  206. chunk = pre.read(512)
  207. if chunk:
  208. flushed += len(chunk)
  209. quiet_for = 0.0
  210. else:
  211. quiet_for += 0.05
  212. time.sleep(0.05)
  213. pre.close()
  214. if flushed:
  215. log(f"flushed {flushed} stale bytes from kernel TTY buffer")
  216. ser = serial.Serial(TTY, BAUD, timeout=0.05)
  217. ser.reset_input_buffer()
  218. ser.reset_output_buffer()
  219. log(f"opened {TTY} @ {BAUD} fd={ser.fd}")
  220. except Exception as e:
  221. log(f"could not open {TTY}: {e}")
  222. sys.exit(1)
  223. def tty_to_tcp(client):
  224. while True:
  225. try:
  226. data = ser.read(256)
  227. if data:
  228. client.sendall(data)
  229. except Exception:
  230. return
  231. srv = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
  232. srv.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
  233. srv.bind(("127.0.0.1", PORT))
  234. srv.listen(5)
  235. log(f"listening on 127.0.0.1:{PORT}")
  236. while True:
  237. client, addr = srv.accept()
  238. log(f"client {addr}")
  239. threading.Thread(target=tty_to_tcp, args=(client,), daemon=True).start()
  240. try:
  241. while True:
  242. data = client.recv(4096)
  243. if not data:
  244. break
  245. ser.write(data)
  246. ser.flush()
  247. except Exception as e:
  248. log(f"tcp err: {e}")
  249. try:
  250. client.close()
  251. except Exception:
  252. pass
  253. log(f"client closed")
  254. '''
  255. class _SSHTunneledSerial(Transport):
  256. """
  257. Auto-managed SSH transport: connects to a remote host, uploads a
  258. tiny pyserial-based TCP<->serial proxy via SFTP, executes it as a
  259. detached background process, and thereby exposes the remote
  260. /dev/ttyUSBX as a local TCP socket reachable through an SSH
  261. direct-tcpip channel.
  262. Cleaning up on close: the proxy process is killed and the uploaded
  263. proxy script is removed.
  264. """
  265. PROXY_REMOTE_PATH = "/tmp/.must_pv1800_proxy.py"
  266. def __init__(self, ssh_host: str, ssh_user: str, ssh_pass: str,
  267. ssh_port: int, serial_port: str, baudrate: int):
  268. if paramiko is None:
  269. raise RuntimeError("paramiko is required for SSH transport; "
  270. "pip install paramiko")
  271. self._client = paramiko.SSHClient()
  272. self._client.set_missing_host_key_policy(paramiko.AutoAddPolicy())
  273. self._client.connect(
  274. hostname=ssh_host, port=ssh_port, username=ssh_user,
  275. password=ssh_pass, allow_agent=False, look_for_keys=False,
  276. timeout=10, banner_timeout=10, auth_timeout=10,
  277. )
  278. transport = self._client.get_transport()
  279. # Make sure no zombies hold /dev/ttyUSBX. Anything left over from
  280. # a previous crash will eat every Modbus response.
  281. ch_k = transport.open_session(timeout=10)
  282. ch_k.exec_command(
  283. f"pkill -9 -f 'socat.*{serial_port}' 2>/dev/null; "
  284. f"pkill -9 -f 'must_pv1800_proxy.py' 2>/dev/null; "
  285. f"pkill -9 -f 'cat {serial_port}' 2>/dev/null; "
  286. f"pkill -9 -f 'dd if={serial_port}' 2>/dev/null; "
  287. f"fuser -k {serial_port} 2>/dev/null; true"
  288. )
  289. ch_k.recv_exit_status()
  290. time.sleep(0.3)
  291. # Make sure pyserial is installed on the remote host.
  292. ch_py = transport.open_session(timeout=30)
  293. ch_py.exec_command("python3 -c 'import serial' 2>/dev/null "
  294. "|| (apt-get install -y python3-serial 2>&1 "
  295. " | tail -3)")
  296. ch_py.recv_exit_status()
  297. try:
  298. ch_py.recv(8192)
  299. except Exception:
  300. pass
  301. # Push the proxy script over SFTP.
  302. sftp = self._client.open_sftp()
  303. try:
  304. with sftp.open(self.PROXY_REMOTE_PATH, "w") as f:
  305. f.write(_REMOTE_PROXY_SCRIPT)
  306. sftp.chmod(self.PROXY_REMOTE_PATH, 0o755)
  307. finally:
  308. sftp.close()
  309. # Pick a free TCP port on the remote side (9500-9600).
  310. self._remote_port = None
  311. for p in range(9500, 9600):
  312. ch_p = transport.open_session(timeout=5)
  313. ch_p.exec_command(
  314. f"ss -tlnH 'sport = :{p}' 2>/dev/null | grep -q LISTEN && "
  315. f"echo taken || echo free"
  316. )
  317. ch_p.recv_exit_status()
  318. if b"free" in ch_p.recv(1024):
  319. self._remote_port = p
  320. break
  321. if self._remote_port is None:
  322. self._client.close()
  323. raise RuntimeError("could not find a free TCP port on remote host")
  324. # Launch the proxy detached from the SSH channel. We use a
  325. # non-looping command so the captured stdout is just the single
  326. # `ss` line we care about.
  327. ch_b = transport.open_session(timeout=15)
  328. ch_b.exec_command(
  329. f"PROXY_PORT={self._remote_port} PROXY_TTY={serial_port} "
  330. f"PROXY_BAUD={baudrate} setsid nohup python3 "
  331. f"{self.PROXY_REMOTE_PATH} </dev/null "
  332. f">/tmp/.must_pv1800_proxy.log 2>&1 & disown; "
  333. f"for i in 1 2 3 4 5 6 7 8 9 10; do "
  334. f" out=$(ss -tlnH 'sport = :{self._remote_port}'); "
  335. f" if [ -n \"$out\" ]; then echo \"$out\"; break; fi; "
  336. f" sleep 0.5; "
  337. f"done"
  338. )
  339. ch_b.settimeout(10)
  340. listener_info = b""
  341. try:
  342. while True:
  343. chunk = ch_b.recv(4096)
  344. if not chunk:
  345. break
  346. listener_info += chunk
  347. except Exception:
  348. pass
  349. try:
  350. ch_b.recv_exit_status()
  351. except Exception:
  352. pass
  353. listener_info = listener_info.decode(errors="replace")
  354. if "LISTEN" not in listener_info:
  355. self._client.close()
  356. raise RuntimeError(
  357. f"proxy did not bind 127.0.0.1:{self._remote_port} on "
  358. f"remote. listener_info={listener_info!r} "
  359. f"log: " + self._safe_remote_log(transport)
  360. )
  361. # Open a direct-tcpip SSH channel from local -> remote 127.0.0.1:port.
  362. self._chan = transport.open_channel(
  363. "direct-tcpip",
  364. ("127.0.0.1", self._remote_port),
  365. ("127.0.0.1", 0),
  366. )
  367. self._chan.settimeout(RESPONSE_TIMEOUT)
  368. self._stash = b""
  369. @staticmethod
  370. def _safe_remote_log(transport) -> str:
  371. try:
  372. ch = transport.open_session(timeout=3)
  373. ch.exec_command("cat /tmp/.must_pv1800_proxy.log 2>&1 | tail -5")
  374. ch.recv_exit_status()
  375. return ch.recv(4096).decode(errors="replace")
  376. except Exception:
  377. return "(log unavailable)"
  378. def drain_input(self) -> None:
  379. """Drain everything pending on the channel, then wait for the
  380. channel to be QUIET for at least 100 ms before returning. This
  381. guarantees that no straggler bytes from a previous oversized
  382. inverter dump are still in flight when we issue the next request.
  383. The 100 ms quiet period was chosen to be larger than the time it
  384. takes the proxy to forward one full Modbus RTU frame at 19200
  385. baud (a 100-byte frame at 8N1 is ~52 ms; 100 ms gives a safe
  386. margin).
  387. """
  388. self._stash = b""
  389. # First, non-blocking drain of anything immediately available.
  390. try:
  391. while True:
  392. r, _, _ = select.select([self._chan], [], [], 0)
  393. if not r:
  394. break
  395. self._chan.recv(4096)
  396. except Exception:
  397. pass
  398. # Then wait for the line to stay quiet for 100 ms.
  399. quiet_for = 0.0
  400. deadline = time.monotonic() + 1.0 # max total drain time
  401. while quiet_for < 0.1 and time.monotonic() < deadline:
  402. r, _, _ = select.select([self._chan], [], [], 0.02)
  403. if r:
  404. try:
  405. self._chan.recv(4096)
  406. except Exception:
  407. break
  408. quiet_for = 0.0
  409. else:
  410. quiet_for += 0.02
  411. def _stash_extra(self, data: bytes) -> None:
  412. self._stash = getattr(self, "_stash", b"") + data
  413. def write(self, data: bytes) -> int:
  414. self._chan.sendall(data)
  415. return len(data)
  416. def read_exact(self, n: int, timeout: float) -> bytes:
  417. # Consume from the stash first (bytes the previous call didn't use).
  418. stash = getattr(self, "_stash", b"")
  419. if stash:
  420. take = stash[:n]
  421. self._stash = stash[n:]
  422. if len(take) >= n:
  423. return bytes(take)
  424. chunks = bytearray(take)
  425. else:
  426. chunks = bytearray()
  427. deadline = time.monotonic() + max(0.05, timeout)
  428. while len(chunks) < n and time.monotonic() < deadline:
  429. remaining = deadline - time.monotonic()
  430. r, _, _ = select.select([self._chan], [], [], remaining)
  431. if not r:
  432. break
  433. try:
  434. chunk = self._chan.recv(n - len(chunks))
  435. if not chunk:
  436. break
  437. chunks.extend(chunk)
  438. except Exception:
  439. break
  440. return bytes(chunks)
  441. def close(self) -> None:
  442. try:
  443. self._chan.close()
  444. except Exception:
  445. pass
  446. try:
  447. transport = self._client.get_transport()
  448. ch = transport.open_session(timeout=5)
  449. ch.exec_command(
  450. f"pkill -9 -f 'must_pv1800_proxy.py' 2>/dev/null; "
  451. f"rm -f {self.PROXY_REMOTE_PATH}; true"
  452. )
  453. ch.recv_exit_status()
  454. except Exception:
  455. pass
  456. try:
  457. self._client.close()
  458. except Exception:
  459. pass
  460. # ---- Modbus over Transport -----------------------------------------------
  461. def modbus_read_holding(t: Transport, slave: int, address: int, count: int = 1,
  462. timeout: float = RESPONSE_TIMEOUT) -> List[int]:
  463. """
  464. Send a Modbus function-3 (Read Holding Registers) request. Returns a
  465. list of unsigned 16-bit register values.
  466. Automatically retries on CRC mismatch or truncated responses -- the
  467. MUST PV1800 occasionally drops a byte at the start of a burst, so a
  468. single retry is enough to recover in practice.
  469. """
  470. last_err = None
  471. for attempt in range(3):
  472. try:
  473. return _modbus_read_holding_once(t, slave, address, count, timeout)
  474. except IOError as e:
  475. last_err = e
  476. # Drain anything that came in for the failed attempt before
  477. # the next try so we don't feed it back into ourselves.
  478. try:
  479. t.drain_input()
  480. except Exception:
  481. pass
  482. time.sleep(0.05 * (attempt + 1))
  483. raise last_err # type: ignore[misc]
  484. def _modbus_read_holding_once(t: Transport, slave: int, address: int,
  485. count: int, timeout: float) -> List[int]:
  486. """Single-shot Modbus read; raises IOError on failure.
  487. The PV1800 firmware doesn't honour our `count` field -- it often
  488. dumps a 39-register block when we asked for one. So instead of
  489. reading a fixed number of bytes we:
  490. 1. Read 5 bytes (the response header).
  491. 2. Parse the bytecount byte to find the actual frame length.
  492. 3. Read the rest of the frame.
  493. 4. Look for the response somewhere in the buffer (multiple frames
  494. may be concatenated if the inverter dumped more than one
  495. block's worth of data) by scanning for a CRC-valid frame.
  496. """
  497. pdu = bytes((slave, 0x03,
  498. (address >> 8) & 0xFF, address & 0xFF,
  499. (count >> 8) & 0xFF, count & 0xFF))
  500. request_frame = pdu + _crc16(pdu)
  501. t.drain_input()
  502. t.write(request_frame)
  503. expected_bc = 2 * count
  504. # Step 1: read at least the 5-byte header so we know the bytecount.
  505. buf = bytearray()
  506. deadline = time.monotonic() + timeout
  507. while len(buf) < 5 and time.monotonic() < deadline:
  508. chunk = t.read_exact(5 - len(buf),
  509. max(0.05, deadline - time.monotonic()))
  510. if chunk:
  511. buf.extend(chunk)
  512. if len(buf) < 5:
  513. raise IOError(f"short read: got {len(buf)} bytes, expected >= 5 "
  514. f"(buf={buf.hex()})")
  515. # Step 2: if the header looks like a valid response, read the rest.
  516. # The inverter's bc may legitimately be > 2*count (oversized dump).
  517. if (buf[0] == slave and buf[1] == 0x03
  518. and 0 <= buf[2] <= 250):
  519. actual_bc = buf[2]
  520. target_len = 5 + actual_bc
  521. while len(buf) < target_len and time.monotonic() < deadline:
  522. chunk = t.read_exact(target_len - len(buf),
  523. max(0.05, deadline - time.monotonic()))
  524. if chunk:
  525. buf.extend(chunk)
  526. # Step 3: scan for a CRC-valid response inside the buffer.
  527. parsed = _find_valid_response(bytes(buf), slave, expected_bc)
  528. if parsed is None:
  529. raise IOError(f"no CRC-valid response in {len(buf)} bytes "
  530. f"(buf={buf.hex()})")
  531. offset, actual_bc, payload, valid_len = parsed
  532. # Step 4: stash any leftover bytes so the next call picks them up.
  533. leftover = bytes(buf[offset + valid_len:])
  534. if leftover:
  535. t._stash_extra(leftover)
  536. return [int.from_bytes(payload[i:i + 2], "big")
  537. for i in range(0, expected_bc, 2)]
  538. def _find_valid_response(buf: bytes, slave: int,
  539. expected_bc: int) -> Optional[Tuple[int, int, bytes, int]]:
  540. """Scan `buf` for a CRC-valid Modbus RTU response from `slave`.
  541. Returns (offset, actual_bc, payload_bytes, total_len_consumed) on
  542. success, None if no valid response is found.
  543. Multiple frames may be concatenated in the buffer if the inverter
  544. sent them back-to-back. We try every byte alignment starting from
  545. offset 0 and pick the FIRST match with a valid CRC. The CRC check
  546. is what tells us we found a real frame boundary; raw header bytes
  547. alone are not enough because the inverter's large dumps can
  548. contain many coincidental `[04 03 bc]` substrings.
  549. """
  550. for offset in range(len(buf) - 4):
  551. if buf[offset] != slave:
  552. continue
  553. func = buf[offset + 1]
  554. if func not in (0x03, 0x83):
  555. continue
  556. bc = buf[offset + 2]
  557. if bc < expected_bc or bc > 250:
  558. continue
  559. total = 5 + bc
  560. if offset + total > len(buf):
  561. continue
  562. frame = buf[offset:offset + total]
  563. if _crc16(frame[:3 + bc]) != frame[3 + bc:5 + bc]:
  564. continue
  565. payload = frame[3:3 + bc]
  566. return (offset, bc, payload, total)
  567. return None
  568. def u16(v: int) -> int:
  569. return v & 0xFFFF
  570. def s16(v: int) -> int:
  571. v &= 0xFFFF
  572. return v - 0x10000 if v & 0x8000 else v
  573. # ---------------------------------------------------------------------------
  574. # Decoded value tables
  575. # ---------------------------------------------------------------------------
  576. # Cross-checked against the official protocol document
  577. # `PH1800 PV1800 EP1800 PV3500 EP3500 RS485 Modbus RTU communication
  578. # Protocol 1.4.15.xlsx` (taHC81/MUST-ESPhome repo; the same document
  579. # shared on Scribd as document 881604771).
  580. CHARGER_WORKSTATE = {
  581. 0: "Initialization mode",
  582. 1: "SelfTest mode",
  583. 2: "Work mode",
  584. 3: "Stop mode",
  585. }
  586. MPPT_STATE = {
  587. 0: "Stop",
  588. 1: "MPPT",
  589. 2: "Current limiting",
  590. }
  591. CHARGING_STATE = {
  592. 0: "Stop",
  593. 1: "Absorb charge",
  594. 2: "Float charge",
  595. 3: "EQ charge",
  596. }
  597. INVERTER_WORKSTATE = {
  598. 0: "PowerOn",
  599. 1: "SelfTest",
  600. 2: "OffGrid",
  601. 3: "GridTie",
  602. 4: "ByPass",
  603. 5: "Stop",
  604. 6: "GridCharging",
  605. }
  606. BATTERY_TYPE = {
  607. 0: "AGM",
  608. 1: "Flooded",
  609. 2: "User defined battery",
  610. 3: "Lithium",
  611. }
  612. ENERGY_USE_MODE = {
  613. 0: "SBU (Solar/battery/utility)",
  614. 1: "SBU (Solar/battery/utility)",
  615. 2: "SUB (Solar/utility/battery)",
  616. 3: "UTI (Utility only)",
  617. 4: "SOL (Solar only)",
  618. }
  619. SOLARUSE_AIM = {
  620. 0: "LBU (Less battery use)",
  621. 1: "LBU (Less battery use)",
  622. 2: "OSO (Only solar output)",
  623. }
  624. # ---------------------------------------------------------------------------
  625. # Register map (from MUST-ESPhome ESPHome config)
  626. # ---------------------------------------------------------------------------
  627. @dataclass(frozen=True)
  628. class Register:
  629. address: int
  630. name: str
  631. scale: float = 1.0
  632. signed: bool = False
  633. decode: Optional[Callable[[int], Any]] = None
  634. unit: str = ""
  635. # CHARGER block
  636. CHARGER_REGS: List[Register] = [
  637. Register(15201, "Charger workstate", decode=lambda v: CHARGER_WORKSTATE.get(v, f"unknown({v})")),
  638. Register(15202, "MPPT state", decode=lambda v: MPPT_STATE.get(v, f"unknown({v})")),
  639. Register(15203, "Charging state", decode=lambda v: CHARGING_STATE.get(v, f"unknown({v})")),
  640. Register(15205, "PV voltage", scale=0.1, unit="V"),
  641. Register(15206, "Battery voltage", scale=0.1, unit="V"),
  642. Register(15207, "Charger current", scale=0.1, unit="A"),
  643. Register(15208, "Charger power", unit="W"),
  644. # 15211 Battery Relay, 15212 PV Relay -- not exposed by reference config
  645. ]
  646. # INVERTER block (read-only telemetry)
  647. INVERTER_REGS: List[Register] = [
  648. Register(25201, "Inverter Work state", decode=lambda v: INVERTER_WORKSTATE.get(v, f"unknown({v})")),
  649. Register(25205, "Battery voltage", scale=0.1, unit="V"),
  650. Register(25206, "Inverter voltage", scale=0.1, unit="V"),
  651. Register(25207, "Grid voltage", scale=0.1, unit="V"),
  652. Register(25213, "Inverter power", signed=True, unit="W"),
  653. Register(25214, "Grid power", signed=True, unit="W"),
  654. Register(25215, "Load power", signed=True, unit="W"),
  655. Register(25216, "System load", unit="%"),
  656. Register(25233, "AC radiator temp", unit="°C"),
  657. Register(25234, "Transformer temp", unit="°C"),
  658. Register(25235, "DC radiator temp", unit="°C"),
  659. Register(25273, "Battery power", signed=True, unit="W"),
  660. Register(25274, "Battery current", signed=True, scale=0.1, unit="A"),
  661. ]
  662. # SETTINGS block (read-only mirror of writable config)
  663. SETTING_REGS: List[Register] = [
  664. Register(10103, "Float voltage", scale=0.1, unit="V"),
  665. Register(10104, "Absorb voltage", scale=0.1, unit="V"),
  666. Register(20118, "Battery stop discharging voltage", scale=0.1, unit="V"),
  667. Register(20119, "Battery stop charging voltage", scale=0.1, unit="V"),
  668. Register(20127, "Battery low voltage", scale=0.1, unit="V"),
  669. Register(20128, "Battery high voltage", scale=0.1, unit="V"),
  670. Register(20132, "Charger current", scale=0.1, unit="A"),
  671. ]
  672. # ---------------------------------------------------------------------------
  673. # Polling / decoding
  674. # ---------------------------------------------------------------------------
  675. def poll_block(t: Transport, slave: int, regs: List[Register],
  676. base_addr: int) -> List[Tuple[str, Any, str]]:
  677. """
  678. Read the entire contiguous register block starting at `base_addr`
  679. in ONE Modbus read, then pluck out the registers we actually care
  680. about from the response. The PV1800 firmware ignores our `count`
  681. anyway and dumps the whole block, so doing this explicitly makes
  682. the request fast AND unambiguous -- we can identify which frame in
  683. a multi-frame buffer is ours because the response contains the
  684. entire block starting at base_addr.
  685. `regs` should be a list of `Register` whose addresses all fall in
  686. the [base_addr, base_addr + count - 1] range. The function returns
  687. (name, decoded_value, unit) tuples for every register in `regs`,
  688. in the same order.
  689. """
  690. if not regs:
  691. return []
  692. block_len = max(r.address for r in regs) - base_addr + 1
  693. raw_regs = modbus_read_holding(t, slave, base_addr, count=block_len)
  694. # raw_regs[0] == base_addr, raw_regs[i] == base_addr + i.
  695. results: List[Tuple[str, Any, str]] = []
  696. for reg in regs:
  697. idx = reg.address - base_addr
  698. raw = raw_regs[idx]
  699. if reg.decode is not None:
  700. value = reg.decode(raw)
  701. elif reg.signed:
  702. value = s16(raw)
  703. else:
  704. value = raw * reg.scale
  705. results.append((reg.name, value, reg.unit))
  706. return results
  707. def poll_individual(t: Transport, slave: int,
  708. regs: List[Register]) -> List[Tuple[str, Any, str]]:
  709. """
  710. Read each register individually with a small inter-register delay.
  711. Used for settings which are scattered across non-contiguous ranges.
  712. """
  713. results: List[Tuple[str, Any, str]] = []
  714. for reg in regs:
  715. raw = modbus_read_holding(t, slave, reg.address, count=1)[0]
  716. if reg.decode is not None:
  717. value = reg.decode(raw)
  718. elif reg.signed:
  719. value = s16(raw)
  720. else:
  721. value = raw * reg.scale
  722. results.append((reg.name, value, reg.unit))
  723. time.sleep(INTER_REGISTER_DELAY)
  724. return results
  725. def print_snapshot(snapshot: List[Tuple[str, Any, str]], title: str) -> None:
  726. bar = "-" * max(15, len(title) + 4)
  727. print(bar)
  728. print(f" {title}")
  729. print(bar)
  730. name_w = max(len(n) for n, _, _ in snapshot)
  731. for name, value, unit in snapshot:
  732. if isinstance(value, float):
  733. vstr = f"{value:.1f}"
  734. else:
  735. vstr = str(value)
  736. suffix = f"{unit}" if unit else ""
  737. print(f" {name:<{name_w}} = {vstr}{suffix}")
  738. print()
  739. # CLI
  740. # ---------------------------------------------------------------------------
  741. def parse_args() -> argparse.Namespace:
  742. p = argparse.ArgumentParser(
  743. description=__doc__,
  744. formatter_class=argparse.RawDescriptionHelpFormatter,
  745. )
  746. g = p.add_mutually_exclusive_group()
  747. g.add_argument("--serial", dest="mode", action="store_const", const="serial",
  748. help="Open --serial-port locally (default if --tcp/--ssh absent).")
  749. g.add_argument("--tcp", dest="mode", action="store_const", const="tcp",
  750. help="Connect to a TCP bridge (socat/ser2net).")
  751. g.add_argument("--ssh", dest="mode", action="store_const", const="ssh",
  752. help="Tunnel over SSH and start a remote pyserial-based proxy (default).")
  753. p.set_defaults(mode="ssh")
  754. p.add_argument("--tcp-host", default=os.environ.get("MUST_TCP_HOST", "127.0.0.1"))
  755. p.add_argument("--tcp-port", type=int,
  756. default=int(os.environ.get("MUST_TCP_PORT", "8502")))
  757. p.add_argument("--ssh-host", default=os.environ.get("MUST_SSH_HOST", DEFAULT_SSH_HOST))
  758. p.add_argument("--ssh-port", type=int,
  759. default=int(os.environ.get("MUST_SSH_PORT", DEFAULT_SSH_PORT)))
  760. p.add_argument("--ssh-user", default=os.environ.get("MUST_SSH_USER", DEFAULT_SSH_USER))
  761. p.add_argument("--ssh-pass", default=os.environ.get("MUST_SSH_PASS", DEFAULT_SSH_PASS))
  762. p.add_argument("--serial-port", default=os.environ.get("MUST_SERIAL_PORT", DEFAULT_SERIAL_PORT))
  763. p.add_argument("--baudrate", type=int,
  764. default=int(os.environ.get("MUST_BAUD", DEFAULT_BAUDRATE)))
  765. p.add_argument("--slave", type=int,
  766. default=int(os.environ.get("MUST_SLAVE", DEFAULT_SLAVE)))
  767. p.add_argument("--interval", type=float,
  768. default=float(os.environ.get("MUST_INTERVAL", DEFAULT_POLL_INTERVAL)),
  769. help="Polling interval in seconds (default: %(default)s)")
  770. p.add_argument("--once", action="store_true",
  771. help="Run a single snapshot and exit")
  772. p.add_argument("--json", action="store_true",
  773. help="Print each snapshot as a JSON object (one line per snapshot).")
  774. return p.parse_args()
  775. def open_transport(args) -> Tuple[Transport, str]:
  776. """
  777. Open the requested transport and return (transport, label).
  778. """
  779. if args.mode == "serial":
  780. return _SerialTransport(args.serial_port, args.baudrate), \
  781. f"serial {args.serial_port} @ {args.baudrate}"
  782. if args.mode == "tcp":
  783. return _SocketTransport(args.tcp_host, args.tcp_port), \
  784. f"tcp {args.tcp_host}:{args.tcp_port}"
  785. # default: ssh
  786. return _SSHTunneledSerial(
  787. ssh_host=args.ssh_host,
  788. ssh_user=args.ssh_user,
  789. ssh_pass=args.ssh_pass,
  790. ssh_port=args.ssh_port,
  791. serial_port=args.serial_port,
  792. baudrate=args.baudrate,
  793. ), f"ssh {args.ssh_user}@{args.ssh_host} -> {args.serial_port}"
  794. def snapshot_to_dicts(charger, inverter, settings) -> dict:
  795. """Combine three (name,value,unit) lists into a single dict for JSON."""
  796. out = {}
  797. for src in (charger, inverter, settings):
  798. for name, value, unit in src:
  799. if isinstance(value, float):
  800. out[name] = {"value": round(value, 3), "unit": unit}
  801. else:
  802. out[name] = {"value": value, "unit": unit}
  803. return out
  804. def main() -> int:
  805. args = parse_args()
  806. print(f"Connecting to MUST PV1800 ({args.mode}): "
  807. f"slave={args.slave}, interval={args.interval}s",
  808. file=sys.stderr)
  809. try:
  810. t, label = open_transport(args)
  811. except Exception as e:
  812. print(f"ERROR: could not open transport: {e}", file=sys.stderr)
  813. return 1
  814. print(f"Transport: {label}", file=sys.stderr)
  815. try:
  816. # Quick connectivity check: read Charger workstate register.
  817. try:
  818. sanity = modbus_read_holding(t, args.slave, 15201, count=1,
  819. timeout=2.0)
  820. print(f"Sanity OK: register 15201 = {sanity[0]} "
  821. f"({CHARGER_WORKSTATE.get(sanity[0], '?')})",
  822. file=sys.stderr)
  823. except Exception as e:
  824. print(f"ERROR: initial Modbus read failed: {e}", file=sys.stderr)
  825. print("Hints:", file=sys.stderr)
  826. print(" - Confirm /dev/ttyUSB0 exists on the remote host.",
  827. file=sys.stderr)
  828. print(" - Confirm baud rate (default 19200, 8N1).",
  829. file=sys.stderr)
  830. print(" - Confirm slave ID (default 4).", file=sys.stderr)
  831. print(" - Another process might already hold the port "
  832. "(check `fuser /dev/ttyUSB0` over SSH).",
  833. file=sys.stderr)
  834. return 2
  835. next_tick = time.monotonic()
  836. attempts = 0
  837. max_once_attempts = 5 # for --once: try this many times before giving up
  838. while True:
  839. now = time.strftime("%Y-%m-%d %H:%M:%S")
  840. try:
  841. # CHARGER block: reads registers 15201..15208 (8 regs).
  842. # We always start at 15201 even if the firmware returns
  843. # more -- the inverter ignores our count and dumps the
  844. # whole block, so this is the unambiguous read.
  845. charger = poll_block(t, args.slave, CHARGER_REGS, base_addr=15201)
  846. # INVERTER block: reads registers 25201..25274 (74 regs).
  847. inverter = poll_block(t, args.slave, INVERTER_REGS, base_addr=25201)
  848. # SETTINGS are split across two non-contiguous ranges:
  849. # 10103..10104 (2 regs) and 20118..20132 (15 regs).
  850. # Polling them as blocks avoids the per-register cross-
  851. # frame mix-up we get from individual reads.
  852. s1 = poll_block(t, args.slave,
  853. [r for r in SETTING_REGS if r.address < 11000],
  854. base_addr=10103)
  855. s2 = poll_block(t, args.slave,
  856. [r for r in SETTING_REGS if r.address >= 11000],
  857. base_addr=20118)
  858. settings = s1 + s2
  859. except Exception as e:
  860. attempts += 1
  861. print(f"[{now}] poll error: {e}", file=sys.stderr)
  862. if args.once and attempts >= max_once_attempts:
  863. return 3
  864. time.sleep(args.interval)
  865. continue
  866. attempts = 0
  867. if args.json:
  868. import json
  869. payload = {"timestamp": now,
  870. "snapshot": snapshot_to_dicts(
  871. charger, inverter, settings)}
  872. print(json.dumps(payload, ensure_ascii=False))
  873. else:
  874. print(f"=== {now} ===")
  875. print_snapshot(charger, "CHARGER")
  876. print_snapshot(inverter, "INVERTER")
  877. print_snapshot(settings, "SETTINGS")
  878. sys.stdout.flush()
  879. if args.once:
  880. return 0
  881. next_tick += args.interval
  882. sleep_for = next_tick - time.monotonic()
  883. if sleep_for > 0:
  884. time.sleep(sleep_for)
  885. else:
  886. next_tick = time.monotonic()
  887. finally:
  888. try:
  889. t.close()
  890. except Exception:
  891. pass
  892. if __name__ == "__main__":
  893. sys.exit(main())