Du kannst nicht mehr als 25 Themen auswählen Themen müssen entweder mit einem Buchstaben oder einer Ziffer beginnen. Sie können Bindestriche („-“) enthalten und bis zu 35 Zeichen lang sein.
 
 
 
 
 
 

475 Zeilen
16 KiB

  1. """不依赖pyserial的Windows原生串口Modbus RTU客户端。"""
  2. from __future__ import annotations
  3. import ctypes
  4. import queue
  5. import threading
  6. import time
  7. import winreg
  8. from ctypes import wintypes
  9. from dataclasses import dataclass
  10. from typing import List, Optional, Tuple
  11. GENERIC_READ = 0x80000000
  12. GENERIC_WRITE = 0x40000000
  13. OPEN_EXISTING = 3
  14. FILE_ATTRIBUTE_NORMAL = 0x80
  15. PURGE_TXABORT = 0x0001
  16. PURGE_RXABORT = 0x0002
  17. PURGE_TXCLEAR = 0x0004
  18. PURGE_RXCLEAR = 0x0008
  19. INVALID_HANDLE_VALUE = ctypes.c_void_p(-1).value
  20. class DCB(ctypes.Structure):
  21. _fields_ = [
  22. ("DCBlength", wintypes.DWORD),
  23. ("BaudRate", wintypes.DWORD),
  24. ("flags", wintypes.DWORD),
  25. ("wReserved", wintypes.WORD),
  26. ("XonLim", wintypes.WORD),
  27. ("XoffLim", wintypes.WORD),
  28. ("ByteSize", wintypes.BYTE),
  29. ("Parity", wintypes.BYTE),
  30. ("StopBits", wintypes.BYTE),
  31. ("XonChar", ctypes.c_char),
  32. ("XoffChar", ctypes.c_char),
  33. ("ErrorChar", ctypes.c_char),
  34. ("EofChar", ctypes.c_char),
  35. ("EvtChar", ctypes.c_char),
  36. ("wReserved1", wintypes.WORD),
  37. ]
  38. class COMMTIMEOUTS(ctypes.Structure):
  39. _fields_ = [
  40. ("ReadIntervalTimeout", wintypes.DWORD),
  41. ("ReadTotalTimeoutMultiplier", wintypes.DWORD),
  42. ("ReadTotalTimeoutConstant", wintypes.DWORD),
  43. ("WriteTotalTimeoutMultiplier", wintypes.DWORD),
  44. ("WriteTotalTimeoutConstant", wintypes.DWORD),
  45. ]
  46. def list_serial_ports() -> List[str]:
  47. """从Windows注册表枚举COM端口,不需要pyserial。"""
  48. ports = set()
  49. path = r"HARDWARE\DEVICEMAP\SERIALCOMM"
  50. try:
  51. with winreg.OpenKey(winreg.HKEY_LOCAL_MACHINE, path) as key:
  52. index = 0
  53. while True:
  54. try:
  55. _, value, _ = winreg.EnumValue(key, index)
  56. ports.add(str(value))
  57. index += 1
  58. except OSError:
  59. break
  60. except OSError:
  61. pass
  62. def sort_key(name: str) -> Tuple[str, int]:
  63. prefix = "".join(ch for ch in name if not ch.isdigit())
  64. digits = "".join(ch for ch in name if ch.isdigit())
  65. return prefix, int(digits or 0)
  66. return sorted(ports, key=sort_key)
  67. class WinSerialPort:
  68. """使用同步Win32 API访问串口;由后台线程调用。"""
  69. def __init__(self) -> None:
  70. self._kernel32 = ctypes.WinDLL("kernel32", use_last_error=True)
  71. self._handle: Optional[int] = None
  72. self._configure_signatures()
  73. def _configure_signatures(self) -> None:
  74. kernel = self._kernel32
  75. kernel.CreateFileW.argtypes = [
  76. wintypes.LPCWSTR,
  77. wintypes.DWORD,
  78. wintypes.DWORD,
  79. wintypes.LPVOID,
  80. wintypes.DWORD,
  81. wintypes.DWORD,
  82. wintypes.HANDLE,
  83. ]
  84. kernel.CreateFileW.restype = wintypes.HANDLE
  85. kernel.CloseHandle.argtypes = [wintypes.HANDLE]
  86. kernel.CloseHandle.restype = wintypes.BOOL
  87. kernel.BuildCommDCBW.argtypes = [wintypes.LPCWSTR, ctypes.POINTER(DCB)]
  88. kernel.BuildCommDCBW.restype = wintypes.BOOL
  89. kernel.SetCommState.argtypes = [wintypes.HANDLE, ctypes.POINTER(DCB)]
  90. kernel.SetCommState.restype = wintypes.BOOL
  91. kernel.SetCommTimeouts.argtypes = [
  92. wintypes.HANDLE,
  93. ctypes.POINTER(COMMTIMEOUTS),
  94. ]
  95. kernel.SetCommTimeouts.restype = wintypes.BOOL
  96. kernel.SetupComm.argtypes = [wintypes.HANDLE, wintypes.DWORD, wintypes.DWORD]
  97. kernel.SetupComm.restype = wintypes.BOOL
  98. kernel.PurgeComm.argtypes = [wintypes.HANDLE, wintypes.DWORD]
  99. kernel.PurgeComm.restype = wintypes.BOOL
  100. kernel.ReadFile.argtypes = [
  101. wintypes.HANDLE,
  102. wintypes.LPVOID,
  103. wintypes.DWORD,
  104. ctypes.POINTER(wintypes.DWORD),
  105. wintypes.LPVOID,
  106. ]
  107. kernel.ReadFile.restype = wintypes.BOOL
  108. kernel.WriteFile.argtypes = [
  109. wintypes.HANDLE,
  110. wintypes.LPCVOID,
  111. wintypes.DWORD,
  112. ctypes.POINTER(wintypes.DWORD),
  113. wintypes.LPVOID,
  114. ]
  115. kernel.WriteFile.restype = wintypes.BOOL
  116. @property
  117. def is_open(self) -> bool:
  118. return self._handle is not None
  119. def open(self, port_name: str, baud_rate: int = 115200) -> None:
  120. self.close()
  121. device_name = rf"\\.\{port_name}"
  122. handle = self._kernel32.CreateFileW(
  123. device_name,
  124. GENERIC_READ | GENERIC_WRITE,
  125. 0,
  126. None,
  127. OPEN_EXISTING,
  128. FILE_ATTRIBUTE_NORMAL,
  129. None,
  130. )
  131. if handle == INVALID_HANDLE_VALUE:
  132. self._raise_last_error(f"打开{port_name}失败")
  133. self._handle = handle
  134. try:
  135. dcb = DCB()
  136. dcb.DCBlength = ctypes.sizeof(DCB)
  137. config = f"baud={baud_rate} parity=e data=8 stop=1"
  138. if not self._kernel32.BuildCommDCBW(config, ctypes.byref(dcb)):
  139. self._raise_last_error("生成串口参数失败")
  140. if not self._kernel32.SetCommState(self._handle, ctypes.byref(dcb)):
  141. self._raise_last_error("设置串口参数失败")
  142. timeouts = COMMTIMEOUTS(
  143. ReadIntervalTimeout=20,
  144. ReadTotalTimeoutMultiplier=0,
  145. ReadTotalTimeoutConstant=20,
  146. WriteTotalTimeoutMultiplier=0,
  147. WriteTotalTimeoutConstant=300,
  148. )
  149. if not self._kernel32.SetCommTimeouts(
  150. self._handle, ctypes.byref(timeouts)
  151. ):
  152. self._raise_last_error("设置串口超时失败")
  153. if not self._kernel32.SetupComm(self._handle, 4096, 4096):
  154. self._raise_last_error("设置串口缓冲区失败")
  155. self.purge()
  156. except Exception:
  157. self.close()
  158. raise
  159. def close(self) -> None:
  160. if self._handle is not None:
  161. self._kernel32.CloseHandle(self._handle)
  162. self._handle = None
  163. def purge(self) -> None:
  164. self._require_open()
  165. flags = PURGE_TXABORT | PURGE_RXABORT | PURGE_TXCLEAR | PURGE_RXCLEAR
  166. if not self._kernel32.PurgeComm(self._handle, flags):
  167. self._raise_last_error("清理串口缓冲区失败")
  168. def write(self, data: bytes) -> None:
  169. self._require_open()
  170. buffer = ctypes.create_string_buffer(data)
  171. written = wintypes.DWORD(0)
  172. if not self._kernel32.WriteFile(
  173. self._handle,
  174. buffer,
  175. len(data),
  176. ctypes.byref(written),
  177. None,
  178. ):
  179. self._raise_last_error("串口发送失败")
  180. if written.value != len(data):
  181. raise OSError(f"串口只发送了{written.value}/{len(data)}字节")
  182. def read(self, maximum: int = 256) -> bytes:
  183. self._require_open()
  184. buffer = ctypes.create_string_buffer(maximum)
  185. received = wintypes.DWORD(0)
  186. if not self._kernel32.ReadFile(
  187. self._handle,
  188. buffer,
  189. maximum,
  190. ctypes.byref(received),
  191. None,
  192. ):
  193. self._raise_last_error("串口接收失败")
  194. return buffer.raw[: received.value]
  195. def _require_open(self) -> None:
  196. if self._handle is None:
  197. raise OSError("串口尚未打开")
  198. @staticmethod
  199. def _raise_last_error(prefix: str) -> None:
  200. code = ctypes.get_last_error()
  201. raise OSError(code, f"{prefix}:{ctypes.FormatError(code).strip()}")
  202. def crc16(data: bytes) -> int:
  203. crc = 0xFFFF
  204. for byte in data:
  205. crc ^= byte
  206. for _ in range(8):
  207. crc = ((crc >> 1) ^ 0xA001) if (crc & 1) else (crc >> 1)
  208. return crc & 0xFFFF
  209. def add_crc(data: bytes) -> bytes:
  210. crc = crc16(data)
  211. return data + bytes((crc & 0xFF, (crc >> 8) & 0xFF))
  212. def valid_crc(frame: bytes) -> bool:
  213. if len(frame) < 4:
  214. return False
  215. received = frame[-2] | (frame[-1] << 8)
  216. return crc16(frame[:-2]) == received
  217. def u16(value: int) -> bytes:
  218. return bytes(((value >> 8) & 0xFF, value & 0xFF))
  219. def read_u16(data: bytes, offset: int) -> int:
  220. return (data[offset] << 8) | data[offset + 1]
  221. def build_read_holding(slave: int, address: int, quantity: int) -> bytes:
  222. if not 1 <= slave <= 247:
  223. raise ValueError("从站地址必须在1到247之间")
  224. if not 1 <= quantity <= 125:
  225. raise ValueError("读取数量必须在1到125之间")
  226. return add_crc(bytes((slave, 0x03)) + u16(address) + u16(quantity))
  227. def build_write_single(slave: int, address: int, value: int) -> bytes:
  228. if not 1 <= slave <= 247:
  229. raise ValueError("从站地址必须在1到247之间")
  230. return add_crc(bytes((slave, 0x06)) + u16(address) + u16(value))
  231. def build_write_multiple(slave: int, address: int, values: List[int]) -> bytes:
  232. if not 1 <= slave <= 247:
  233. raise ValueError("从站地址必须在1到247之间")
  234. if not 1 <= len(values) <= 123:
  235. raise ValueError("写入数量必须在1到123之间")
  236. payload = b"".join(u16(value) for value in values)
  237. body = (
  238. bytes((slave, 0x10))
  239. + u16(address)
  240. + u16(len(values))
  241. + bytes((len(payload),))
  242. + payload
  243. )
  244. return add_crc(body)
  245. @dataclass(frozen=True)
  246. class Request:
  247. frame: bytes
  248. function: int
  249. address: int
  250. quantity_or_value: int
  251. context: str
  252. class ModbusClient:
  253. """单事务后台客户端;GUI通过events队列接收结果。"""
  254. # STM32从站由RTOS任务轮询。连续请求之间留出恢复时间,避免上一帧刚发送
  255. # 完成时下一帧已经进入USART,导致从站状态机漏掉请求。
  256. INTER_REQUEST_DELAY_S = 0.03
  257. RESPONSE_TIMEOUT_S = 0.8
  258. MAX_ATTEMPTS = 2
  259. def __init__(self) -> None:
  260. self.events: "queue.Queue[tuple]" = queue.Queue()
  261. self._requests: "queue.Queue[Optional[Request]]" = queue.Queue()
  262. self._serial = WinSerialPort()
  263. self._thread: Optional[threading.Thread] = None
  264. self._stop_event = threading.Event()
  265. self._busy_event = threading.Event()
  266. self.slave_address = 1
  267. self.port_name = ""
  268. self._last_request_finished = 0.0
  269. @property
  270. def is_open(self) -> bool:
  271. return self._serial.is_open
  272. @property
  273. def is_busy(self) -> bool:
  274. return self._busy_event.is_set() or not self._requests.empty()
  275. def open(self, port_name: str, baud_rate: int, slave_address: int) -> None:
  276. if not 1 <= slave_address <= 247:
  277. raise ValueError("从站地址必须在1到247之间")
  278. self.close()
  279. self._serial.open(port_name, baud_rate)
  280. self.slave_address = slave_address
  281. self.port_name = port_name
  282. self._stop_event.clear()
  283. self._thread = threading.Thread(
  284. target=self._worker_loop,
  285. name="ModbusRtuWorker",
  286. daemon=True,
  287. )
  288. self._thread.start()
  289. def close(self) -> None:
  290. self._stop_event.set()
  291. if self._thread is not None and self._thread.is_alive():
  292. self._requests.put(None)
  293. self._thread.join(timeout=1.0)
  294. self._thread = None
  295. self._busy_event.clear()
  296. while not self._requests.empty():
  297. try:
  298. self._requests.get_nowait()
  299. except queue.Empty:
  300. break
  301. self._serial.close()
  302. self._last_request_finished = 0.0
  303. def read_holding(self, address: int, quantity: int, context: str = "") -> None:
  304. frame = build_read_holding(self.slave_address, address, quantity)
  305. self._requests.put(Request(frame, 0x03, address, quantity, context))
  306. def write_single(self, address: int, value: int, context: str = "") -> None:
  307. frame = build_write_single(self.slave_address, address, value)
  308. self._requests.put(Request(frame, 0x06, address, value, context))
  309. def write_multiple(
  310. self, address: int, values: List[int], context: str = ""
  311. ) -> None:
  312. frame = build_write_multiple(self.slave_address, address, values)
  313. self._requests.put(Request(frame, 0x10, address, len(values), context))
  314. def _worker_loop(self) -> None:
  315. while not self._stop_event.is_set():
  316. try:
  317. request = self._requests.get(timeout=0.1)
  318. except queue.Empty:
  319. continue
  320. if request is None:
  321. break
  322. self._busy_event.set()
  323. try:
  324. for attempt in range(1, self.MAX_ATTEMPTS + 1):
  325. try:
  326. self._wait_inter_request_gap()
  327. self._execute(request)
  328. break
  329. except TimeoutError:
  330. if attempt >= self.MAX_ATTEMPTS:
  331. raise
  332. self.events.put(("retry", request.context, attempt + 1))
  333. except Exception as exc: # 将后台异常送回GUI线程
  334. self.events.put(("error", request.context, str(exc)))
  335. finally:
  336. self._last_request_finished = time.monotonic()
  337. self._busy_event.clear()
  338. def _wait_inter_request_gap(self) -> None:
  339. remaining = (
  340. self._last_request_finished + self.INTER_REQUEST_DELAY_S
  341. - time.monotonic()
  342. )
  343. if remaining > 0:
  344. time.sleep(remaining)
  345. def _execute(self, request: Request) -> None:
  346. self._serial.purge()
  347. self._serial.write(request.frame)
  348. self.events.put(("tx", request.frame))
  349. deadline = time.monotonic() + self.RESPONSE_TIMEOUT_S
  350. response = bytearray()
  351. expected_length: Optional[int] = None
  352. while time.monotonic() < deadline and not self._stop_event.is_set():
  353. chunk = self._serial.read(256)
  354. if chunk:
  355. response.extend(chunk)
  356. if len(response) >= 2:
  357. function = response[1]
  358. if function == (request.function | 0x80):
  359. expected_length = 5
  360. elif function == 0x03 and len(response) >= 3:
  361. expected_length = 5 + response[2]
  362. elif function in (0x06, 0x10):
  363. expected_length = 8
  364. if expected_length is not None and len(response) >= expected_length:
  365. break
  366. if expected_length is None or len(response) < expected_length:
  367. raise TimeoutError("等待从站响应超时")
  368. frame = bytes(response[:expected_length])
  369. self.events.put(("rx", frame))
  370. if not valid_crc(frame):
  371. raise ValueError("响应CRC校验失败")
  372. if frame[0] != self.slave_address:
  373. raise ValueError("响应从站地址不匹配")
  374. function = frame[1]
  375. if function == (request.function | 0x80):
  376. self.events.put(
  377. ("exception", request.function, frame[2], request.context)
  378. )
  379. return
  380. if function != request.function:
  381. raise ValueError("响应功能码不匹配")
  382. if function == 0x03:
  383. byte_count = frame[2]
  384. if byte_count != request.quantity_or_value * 2:
  385. raise ValueError("读取响应数据长度不匹配")
  386. values = [
  387. read_u16(frame, offset)
  388. for offset in range(3, 3 + byte_count, 2)
  389. ]
  390. self.events.put(
  391. ("read", request.address, values, request.context)
  392. )
  393. return
  394. echoed_address = read_u16(frame, 2)
  395. echoed_value = read_u16(frame, 4)
  396. if (
  397. echoed_address != request.address
  398. or echoed_value != request.quantity_or_value
  399. ):
  400. raise ValueError("写响应回显与请求不一致")
  401. self.events.put(
  402. (
  403. "write",
  404. function,
  405. echoed_address,
  406. echoed_value,
  407. request.context,
  408. )
  409. )