must_pv1800_monitor.py 47 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197
  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. def RELAY_STATE(v: int) -> str:
  607. """Decode relay-state registers (25237..25242, 15211, 15212).
  608. Per the protocol: 0=Disconnect, 1=Connect.
  609. """
  610. return "Closed" if v == 1 else "Open"
  611. BATTERY_TYPE = {
  612. 0: "AGM",
  613. 1: "Flooded",
  614. 2: "User defined battery",
  615. 3: "Lithium",
  616. }
  617. ENERGY_USE_MODE = {
  618. 0: "SBU (Solar/battery/utility)",
  619. 1: "SBU (Solar/battery/utility)",
  620. 2: "SUB (Solar/utility/battery)",
  621. 3: "UTI (Utility only)",
  622. 4: "SOL (Solar only)",
  623. }
  624. SOLARUSE_AIM = {
  625. 0: "LBU (Less battery use)",
  626. 1: "LBU (Less battery use)",
  627. 2: "OSO (Only solar output)",
  628. }
  629. # ---------------------------------------------------------------------------
  630. # Register map (from MUST-ESPhome ESPHome config)
  631. # ---------------------------------------------------------------------------
  632. @dataclass(frozen=True)
  633. class Register:
  634. address: int
  635. name: str
  636. scale: float = 1.0
  637. signed: bool = False
  638. decode: Optional[Callable[[int], Any]] = None
  639. unit: str = ""
  640. # CHARGER block (registers 15201..15221; gaps at 15204, 15213, 15214
  641. # are reserved / unused by the firmware)
  642. CHARGER_REGS: List[Register] = [
  643. # ----- Operating state -----
  644. Register(15201, "Charger workstate", decode=lambda v: CHARGER_WORKSTATE.get(v, f"unknown({v})")),
  645. Register(15202, "MPPT state", decode=lambda v: MPPT_STATE.get(v, f"unknown({v})")),
  646. Register(15203, "Charging state", decode=lambda v: CHARGING_STATE.get(v, f"unknown({v})")),
  647. # ----- Live measurements -----
  648. Register(15205, "PV voltage", scale=0.1, unit="V"),
  649. Register(15206, "Battery voltage (charger)", scale=0.1, unit="V"),
  650. Register(15207, "Charger current", scale=0.1, unit="A"),
  651. Register(15208, "Charger power", unit="W"),
  652. # ----- Sensors -----
  653. Register(15209, "Charger radiator temp", unit="°C"),
  654. Register(15210, "External temp", unit="°C"),
  655. # ----- Relays -----
  656. Register(15211, "Battery Relay (charger)", decode=RELAY_STATE),
  657. Register(15212, "PV Relay (charger)", decode=RELAY_STATE),
  658. # ----- Configuration / identity -----
  659. Register(15215, "Battery voltage grade", decode=lambda v: f"{v}V system"),
  660. Register(15216, "Charger rated current", scale=0.1, unit="A"),
  661. # ----- Energy accumulator (high * 1000 kWh + low * 0.1 kWh) -----
  662. Register(15217, "Acc PV energy MWh part", unit="MWh part"),
  663. Register(15218, "Acc PV energy kWh part", scale=0.1, unit="kWh part"),
  664. # ----- Runtime since last reset -----
  665. Register(15219, "Acc runtime days", unit="d"),
  666. Register(15220, "Acc runtime hours", unit="h"),
  667. Register(15221, "Acc runtime minutes", unit="m"),
  668. ]
  669. # INVERTER block (registers 25201..25274 with gaps at 25204, 25220,
  670. # 25224, 25227, 25228, 25232, 25236, 25243, 25244, 25262..25268, 25270..25272,
  671. # 25275, 25276).
  672. INVERTER_REGS: List[Register] = [
  673. # ----- Operating state -----
  674. Register(25201, "Inverter Work state", decode=lambda v: INVERTER_WORKSTATE.get(v, f"unknown({v})")),
  675. Register(25202, "AC voltage grade", decode=lambda v: f"{v}V variant"),
  676. Register(25203, "Rated power", unit="VA"),
  677. # ----- Voltages -----
  678. Register(25205, "Battery voltage (inverter)", scale=0.1, unit="V"),
  679. Register(25206, "Inverter voltage", scale=0.1, unit="V"),
  680. Register(25207, "Grid voltage", scale=0.1, unit="V"),
  681. Register(25208, "BUS voltage", scale=0.1, unit="V"),
  682. # ----- Currents -----
  683. Register(25209, "Control current", scale=0.1, unit="A"),
  684. Register(25210, "Inverter current", scale=0.1, unit="A"),
  685. Register(25211, "Grid current", scale=0.1, unit="A"),
  686. Register(25212, "Load current", scale=0.1, unit="A"),
  687. # ----- Real power (P) -----
  688. Register(25213, "PInverter", signed=True, unit="W"),
  689. Register(25214, "PGrid", signed=True, unit="W"),
  690. Register(25215, "PLoad", signed=True, unit="W"),
  691. Register(25216, "Load percent", unit="%"),
  692. # ----- Apparent power (S) -----
  693. Register(25217, "SInverter", signed=True, unit="VA"),
  694. Register(25218, "SGrid", signed=True, unit="VA"),
  695. Register(25219, "SLoad", signed=True, unit="VA"),
  696. # ----- Reactive power (Q) -----
  697. Register(25221, "QInverter", signed=True, unit="var"),
  698. Register(25222, "QGrid", signed=True, unit="var"),
  699. Register(25223, "QLoad", signed=True, unit="var"),
  700. # ----- Frequencies -----
  701. Register(25225, "Inverter frequency", scale=0.01, unit="Hz"),
  702. Register(25226, "Grid frequency", scale=0.01, unit="Hz"),
  703. # ----- Temperatures -----
  704. Register(25233, "AC radiator temp", unit="°C"),
  705. Register(25234, "Transformer temp", unit="°C"),
  706. Register(25235, "DC radiator temp", unit="°C"),
  707. # ----- Relays -----
  708. Register(25237, "Inverter relay", decode=RELAY_STATE),
  709. Register(25238, "Grid relay", decode=RELAY_STATE),
  710. Register(25239, "Load relay", decode=RELAY_STATE),
  711. Register(25240, "N_Line relay", decode=RELAY_STATE),
  712. Register(25241, "DC relay", decode=RELAY_STATE),
  713. Register(25242, "Earth relay", decode=RELAY_STATE),
  714. # ----- Energy accumulators (each = high*1000 kWh + low*0.1 kWh) -----
  715. Register(25245, "Acc charger MWh part", unit="MWh part"),
  716. Register(25246, "Acc charger kWh part", scale=0.1, unit="kWh part"),
  717. Register(25247, "Acc discharger MWh part", unit="MWh part"),
  718. Register(25248, "Acc discharger kWh part", scale=0.1, unit="kWh part"),
  719. Register(25249, "Acc grid-buy MWh part", unit="MWh part"),
  720. Register(25250, "Acc grid-buy kWh part", scale=0.1, unit="kWh part"),
  721. Register(25251, "Acc grid-sell MWh part", unit="MWh part"),
  722. Register(25252, "Acc grid-sell kWh part", scale=0.1, unit="kWh part"),
  723. Register(25253, "Acc load MWh part", unit="MWh part"),
  724. Register(25254, "Acc load kWh part", scale=0.1, unit="kWh part"),
  725. Register(25255, "Acc self-use MWh part", unit="MWh part"),
  726. Register(25256, "Acc self-use kWh part", scale=0.1, unit="kWh part"),
  727. Register(25257, "Acc PV-sell MWh part", unit="MWh part"),
  728. Register(25258, "Acc PV-sell kWh part", scale=0.1, unit="kWh part"),
  729. Register(25259, "Acc grid-charge MWh part", unit="MWh part"),
  730. Register(25260, "Acc grid-charge kWh part", scale=0.1, unit="kWh part"),
  731. # ----- Battery telemetry (live) -----
  732. Register(25273, "Battery power", signed=True, unit="W"),
  733. Register(25274, "Battery current", signed=True, scale=0.1, unit="A"),
  734. Register(25275, "Battery voltage grade", decode=lambda v: f"{v}V system"),
  735. ]
  736. # SETTINGS block (read-only mirror of writable config)
  737. SETTING_REGS: List[Register] = [
  738. Register(10103, "Float voltage", scale=0.1, unit="V"),
  739. Register(10104, "Absorb voltage", scale=0.1, unit="V"),
  740. Register(20118, "Battery stop discharging voltage", scale=0.1, unit="V"),
  741. Register(20119, "Battery stop charging voltage", scale=0.1, unit="V"),
  742. Register(20127, "Battery low voltage", scale=0.1, unit="V"),
  743. Register(20128, "Battery high voltage", scale=0.1, unit="V"),
  744. Register(20132, "Charger current", scale=0.1, unit="A"),
  745. ]
  746. # ---------------------------------------------------------------------------
  747. # Polling / decoding
  748. # ---------------------------------------------------------------------------
  749. def poll_block(t: Transport, slave: int, regs: List[Register],
  750. base_addr: int) -> List[Tuple[str, Any, str]]:
  751. """
  752. Read the entire contiguous register block starting at `base_addr`
  753. in ONE Modbus read, then pluck out the registers we actually care
  754. about from the response. The PV1800 firmware ignores our `count`
  755. anyway and dumps the whole block, so doing this explicitly makes
  756. the request fast AND unambiguous -- we can identify which frame in
  757. a multi-frame buffer is ours because the response contains the
  758. entire block starting at base_addr.
  759. `regs` should be a list of `Register` whose addresses all fall in
  760. the [base_addr, base_addr + count - 1] range. The function returns
  761. (name, decoded_value, unit) tuples for every register in `regs`,
  762. in the same order.
  763. """
  764. if not regs:
  765. return []
  766. block_len = max(r.address for r in regs) - base_addr + 1
  767. raw_regs = modbus_read_holding(t, slave, base_addr, count=block_len)
  768. # raw_regs[0] == base_addr, raw_regs[i] == base_addr + i.
  769. results: List[Tuple[str, Any, str]] = []
  770. for reg in regs:
  771. idx = reg.address - base_addr
  772. raw = raw_regs[idx]
  773. if reg.decode is not None:
  774. value = reg.decode(raw)
  775. elif reg.signed:
  776. value = s16(raw)
  777. else:
  778. value = raw * reg.scale
  779. results.append((reg.name, value, reg.unit))
  780. return results
  781. def poll_individual(t: Transport, slave: int,
  782. regs: List[Register]) -> List[Tuple[str, Any, str]]:
  783. """
  784. Read each register individually with a small inter-register delay.
  785. Used for settings which are scattered across non-contiguous ranges.
  786. """
  787. results: List[Tuple[str, Any, str]] = []
  788. for reg in regs:
  789. raw = modbus_read_holding(t, slave, reg.address, count=1)[0]
  790. if reg.decode is not None:
  791. value = reg.decode(raw)
  792. elif reg.signed:
  793. value = s16(raw)
  794. else:
  795. value = raw * reg.scale
  796. results.append((reg.name, value, reg.unit))
  797. time.sleep(INTER_REGISTER_DELAY)
  798. return results
  799. def _split_groups(snapshot: List[Tuple[str, Any, str]],
  800. headings: List[Tuple[str, Any]]) -> List[Tuple[str, List[Tuple[str, Any, str]]]]:
  801. """Split a flat snapshot into labelled groups.
  802. `headings` is a list of `(heading_text, name_prefix_or_prefixes)`
  803. pairs. Each register whose name starts with one of the listed
  804. prefixes goes into that group. A prefix can also be a `str` for
  805. the common single-prefix case. Registers that don't match any
  806. prefix are silently dropped.
  807. """
  808. out: List[Tuple[str, List[Tuple[str, Any, str]]]] = []
  809. for header, prefixes in headings:
  810. if isinstance(prefixes, str):
  811. prefixes = [prefixes]
  812. items = [(n, v, u) for n, v, u in snapshot
  813. if any(n.startswith(p) for p in prefixes)]
  814. if items:
  815. out.append((header, items))
  816. return out
  817. def _combine_accumulators(snapshot: List[Tuple[str, Any, str]]) -> List[Tuple[str, Any, str]]:
  818. """Merge every 'Acc <name> MWh part' / 'kWh part' pair into a single
  819. 'Acc <name> (total) = X.X kWh' line. Returns the merged list (in
  820. the original order, with the combined line replacing the two
  821. originals)."""
  822. by_name = {n: (v, u) for n, v, u in snapshot}
  823. out: List[Tuple[str, Any, str]] = []
  824. skip = set()
  825. for name, value, unit in snapshot:
  826. if name in skip:
  827. continue
  828. if name.endswith(" MWh part"):
  829. base = name[: -len(" MWh part")]
  830. kwh_name = base + " kWh part"
  831. if kwh_name in by_name:
  832. kwh_val = by_name[kwh_name][0]
  833. total_kwh = value * 1000.0 + kwh_val * 0.1
  834. # `base` looks like "Acc <name>"; we want "<name>".
  835. if base.startswith("Acc "):
  836. label = base[len("Acc "):]
  837. else:
  838. label = base
  839. out.append((f"Acc {label.strip()} (total)",
  840. round(total_kwh, 2), "kWh"))
  841. skip.add(name)
  842. skip.add(kwh_name)
  843. continue
  844. out.append((name, value, unit))
  845. return out
  846. def print_snapshot(snapshot: List[Tuple[str, Any, str]], title: str,
  847. groups: Optional[List[Tuple[str, List[Tuple[str, Any, str]]]]] = None) -> None:
  848. """Print a snapshot. If `groups` is given, print section headers between
  849. groups; otherwise dump the whole list flat."""
  850. bar = "-" * max(15, len(title) + 4)
  851. print(bar)
  852. print(f" {title}")
  853. print(bar)
  854. name_w = max(len(n) for n, _, _ in snapshot)
  855. if groups:
  856. for header, items in groups:
  857. print(f" --- {header} ---")
  858. for name, value, unit in items:
  859. if isinstance(value, float):
  860. vstr = f"{value:.1f}"
  861. else:
  862. vstr = str(value)
  863. suffix = f"{unit}" if unit else ""
  864. print(f" {name:<{name_w}} = {vstr}{suffix}")
  865. print()
  866. else:
  867. for name, value, unit in snapshot:
  868. if isinstance(value, float):
  869. vstr = f"{value:.1f}"
  870. else:
  871. vstr = str(value)
  872. suffix = f"{unit}" if unit else ""
  873. print(f" {name:<{name_w}} = {vstr}{suffix}")
  874. print()
  875. # CLI
  876. # ---------------------------------------------------------------------------
  877. def parse_args() -> argparse.Namespace:
  878. p = argparse.ArgumentParser(
  879. description=__doc__,
  880. formatter_class=argparse.RawDescriptionHelpFormatter,
  881. )
  882. g = p.add_mutually_exclusive_group()
  883. g.add_argument("--serial", dest="mode", action="store_const", const="serial",
  884. help="Open --serial-port locally (default if --tcp/--ssh absent).")
  885. g.add_argument("--tcp", dest="mode", action="store_const", const="tcp",
  886. help="Connect to a TCP bridge (socat/ser2net).")
  887. g.add_argument("--ssh", dest="mode", action="store_const", const="ssh",
  888. help="Tunnel over SSH and start a remote pyserial-based proxy (default).")
  889. p.set_defaults(mode="ssh")
  890. p.add_argument("--tcp-host", default=os.environ.get("MUST_TCP_HOST", "127.0.0.1"))
  891. p.add_argument("--tcp-port", type=int,
  892. default=int(os.environ.get("MUST_TCP_PORT", "8502")))
  893. p.add_argument("--ssh-host", default=os.environ.get("MUST_SSH_HOST", DEFAULT_SSH_HOST))
  894. p.add_argument("--ssh-port", type=int,
  895. default=int(os.environ.get("MUST_SSH_PORT", DEFAULT_SSH_PORT)))
  896. p.add_argument("--ssh-user", default=os.environ.get("MUST_SSH_USER", DEFAULT_SSH_USER))
  897. p.add_argument("--ssh-pass", default=os.environ.get("MUST_SSH_PASS", DEFAULT_SSH_PASS))
  898. p.add_argument("--serial-port", default=os.environ.get("MUST_SERIAL_PORT", DEFAULT_SERIAL_PORT))
  899. p.add_argument("--baudrate", type=int,
  900. default=int(os.environ.get("MUST_BAUD", DEFAULT_BAUDRATE)))
  901. p.add_argument("--slave", type=int,
  902. default=int(os.environ.get("MUST_SLAVE", DEFAULT_SLAVE)))
  903. p.add_argument("--interval", type=float,
  904. default=float(os.environ.get("MUST_INTERVAL", DEFAULT_POLL_INTERVAL)),
  905. help="Polling interval in seconds (default: %(default)s)")
  906. p.add_argument("--once", action="store_true",
  907. help="Run a single snapshot and exit")
  908. p.add_argument("--json", action="store_true",
  909. help="Print each snapshot as a JSON object (one line per snapshot).")
  910. return p.parse_args()
  911. def open_transport(args) -> Tuple[Transport, str]:
  912. """
  913. Open the requested transport and return (transport, label).
  914. """
  915. if args.mode == "serial":
  916. return _SerialTransport(args.serial_port, args.baudrate), \
  917. f"serial {args.serial_port} @ {args.baudrate}"
  918. if args.mode == "tcp":
  919. return _SocketTransport(args.tcp_host, args.tcp_port), \
  920. f"tcp {args.tcp_host}:{args.tcp_port}"
  921. # default: ssh
  922. return _SSHTunneledSerial(
  923. ssh_host=args.ssh_host,
  924. ssh_user=args.ssh_user,
  925. ssh_pass=args.ssh_pass,
  926. ssh_port=args.ssh_port,
  927. serial_port=args.serial_port,
  928. baudrate=args.baudrate,
  929. ), f"ssh {args.ssh_user}@{args.ssh_host} -> {args.serial_port}"
  930. def snapshot_to_dicts(charger, inverter, settings) -> dict:
  931. """Combine three (name,value,unit) lists into a single dict for JSON."""
  932. out = {}
  933. for src in (charger, inverter, settings):
  934. for name, value, unit in src:
  935. if isinstance(value, float):
  936. out[name] = {"value": round(value, 3), "unit": unit}
  937. else:
  938. out[name] = {"value": value, "unit": unit}
  939. return out
  940. def main() -> int:
  941. args = parse_args()
  942. print(f"Connecting to MUST PV1800 ({args.mode}): "
  943. f"slave={args.slave}, interval={args.interval}s",
  944. file=sys.stderr)
  945. try:
  946. t, label = open_transport(args)
  947. except Exception as e:
  948. print(f"ERROR: could not open transport: {e}", file=sys.stderr)
  949. return 1
  950. print(f"Transport: {label}", file=sys.stderr)
  951. try:
  952. # Quick connectivity check: read Charger workstate register.
  953. try:
  954. sanity = modbus_read_holding(t, args.slave, 15201, count=1,
  955. timeout=2.0)
  956. print(f"Sanity OK: register 15201 = {sanity[0]} "
  957. f"({CHARGER_WORKSTATE.get(sanity[0], '?')})",
  958. file=sys.stderr)
  959. except Exception as e:
  960. print(f"ERROR: initial Modbus read failed: {e}", file=sys.stderr)
  961. print("Hints:", file=sys.stderr)
  962. print(" - Confirm /dev/ttyUSB0 exists on the remote host.",
  963. file=sys.stderr)
  964. print(" - Confirm baud rate (default 19200, 8N1).",
  965. file=sys.stderr)
  966. print(" - Confirm slave ID (default 4).", file=sys.stderr)
  967. print(" - Another process might already hold the port "
  968. "(check `fuser /dev/ttyUSB0` over SSH).",
  969. file=sys.stderr)
  970. return 2
  971. next_tick = time.monotonic()
  972. attempts = 0
  973. max_once_attempts = 5 # for --once: try this many times before giving up
  974. while True:
  975. now = time.strftime("%Y-%m-%d %H:%M:%S")
  976. try:
  977. # CHARGER block: reads registers 15201..15208 (8 regs).
  978. # We always start at 15201 even if the firmware returns
  979. # more -- the inverter ignores our count and dumps the
  980. # whole block, so this is the unambiguous read.
  981. charger = poll_block(t, args.slave, CHARGER_REGS, base_addr=15201)
  982. # INVERTER block: reads registers 25201..25274 (74 regs).
  983. inverter = poll_block(t, args.slave, INVERTER_REGS, base_addr=25201)
  984. # SETTINGS are split across two non-contiguous ranges:
  985. # 10103..10104 (2 regs) and 20118..20132 (15 regs).
  986. # Polling them as blocks avoids the per-register cross-
  987. # frame mix-up we get from individual reads.
  988. s1 = poll_block(t, args.slave,
  989. [r for r in SETTING_REGS if r.address < 11000],
  990. base_addr=10103)
  991. s2 = poll_block(t, args.slave,
  992. [r for r in SETTING_REGS if r.address >= 11000],
  993. base_addr=20118)
  994. settings = s1 + s2
  995. except Exception as e:
  996. attempts += 1
  997. print(f"[{now}] poll error: {e}", file=sys.stderr)
  998. if args.once and attempts >= max_once_attempts:
  999. return 3
  1000. time.sleep(args.interval)
  1001. continue
  1002. attempts = 0
  1003. if args.json:
  1004. import json
  1005. payload = {"timestamp": now,
  1006. "snapshot": snapshot_to_dicts(
  1007. charger, inverter, settings)}
  1008. print(json.dumps(payload, ensure_ascii=False))
  1009. else:
  1010. print(f"=== {now} ===")
  1011. # Combine accumulator pairs (MWh part + kWh part -> total kWh)
  1012. # and split CHARGER/INVERTER into readable sub-sections.
  1013. charger = _combine_accumulators(charger)
  1014. inverter = _combine_accumulators(inverter)
  1015. charger_groups = _split_groups(charger, [
  1016. ("State", ["Charger workstate", "MPPT state", "Charging state"]),
  1017. ("Live readings", ["PV voltage", "Battery voltage (charger)", "Charger current", "Charger power"]),
  1018. ("Sensors", ["Charger radiator temp", "External temp"]),
  1019. ("Relays", ["Battery Relay (charger)", "PV Relay (charger)"]),
  1020. ("Configuration", ["Battery voltage grade", "Charger rated current"]),
  1021. ("Energy totals", ["Acc PV energy"]),
  1022. ("Runtime", ["Acc runtime "]),
  1023. ])
  1024. inverter_groups = _split_groups(inverter, [
  1025. ("State", ["Inverter Work state"]),
  1026. ("Configuration", ["AC voltage grade", "Rated power", "Battery voltage grade"]),
  1027. ("Voltages", ["Battery voltage (inverter)", "Inverter voltage", "Grid voltage", "BUS voltage"]),
  1028. ("Currents", ["Control current", "Inverter current", "Grid current", "Load current"]),
  1029. ("Real power (P)", ["PInverter", "PGrid", "PLoad", "Load percent"]),
  1030. ("Apparent power (S)", ["SInverter", "SGrid", "SLoad"]),
  1031. ("Reactive power (Q)", ["QInverter", "QGrid", "QLoad"]),
  1032. ("Frequencies", ["Inverter frequency", "Grid frequency"]),
  1033. ("Temperatures", ["AC radiator temp", "Transformer temp", "DC radiator temp"]),
  1034. ("Relays", ["Inverter relay", "Grid relay", "Load relay", "N_Line relay", "DC relay", "Earth relay"]),
  1035. ("Energy totals", ["Acc "]),
  1036. ("Battery telemetry", ["Battery power", "Battery current"]),
  1037. ])
  1038. print_snapshot(charger, "CHARGER", groups=charger_groups)
  1039. print_snapshot(inverter, "INVERTER", groups=inverter_groups)
  1040. print_snapshot(settings, "SETTINGS")
  1041. sys.stdout.flush()
  1042. if args.once:
  1043. return 0
  1044. next_tick += args.interval
  1045. sleep_for = next_tick - time.monotonic()
  1046. if sleep_for > 0:
  1047. time.sleep(sleep_for)
  1048. else:
  1049. next_tick = time.monotonic()
  1050. finally:
  1051. try:
  1052. t.close()
  1053. except Exception:
  1054. pass
  1055. if __name__ == "__main__":
  1056. sys.exit(main())