serial_port.py 27 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705
  1. import serial
  2. import serial.tools.list_ports
  3. import threading
  4. import time
  5. import logging
  6. import platform
  7. import glob
  8. import queue
  9. from dataclasses import dataclass
  10. # 配置日志
  11. logging.basicConfig(
  12. level=logging.INFO,
  13. format='%(asctime)s - %(name)s - %(levelname)s - %(message)s'
  14. )
  15. logger = logging.getLogger('serial_port')
  16. @dataclass
  17. class SerialConfig:
  18. """串口配置数据类"""
  19. port: str
  20. baudrate: int = 115200
  21. bytesize: int = serial.EIGHTBITS
  22. parity: str = serial.PARITY_NONE
  23. stopbits: int = serial.STOPBITS_ONE
  24. timeout: float = 1.0
  25. xonxoff: bool = False
  26. rtscts: bool = False
  27. dsrdtr: bool = False
  28. class SerialPort:
  29. """串口通信类,提供串口连接、读写和状态管理功能"""
  30. def __init__(self):
  31. self.ser = None
  32. self.is_connected = False
  33. self.lock = threading.RLock() # 使用可重入锁
  34. self.read_thread = None
  35. self.stop_event = threading.Event()
  36. self.data_callback = None
  37. self.send_callback = None
  38. self.status_callback = None
  39. self.error_callback = None
  40. self.current_config = None
  41. self.reconnect_attempts = 0
  42. self.max_reconnect_attempts = 3
  43. self.raw_response_buffer = []
  44. self._reconnect_monitor_thread = None
  45. # 同步接收模式:send_and_wait 活动时暂停后台读取线程,防止响应被读线程抢走
  46. self._sync_receive_active = False
  47. self._sync_lock = threading.Lock()
  48. # 串口发送队列:所有发送指令统一排队,按最小间隔逐个发送,避免总线冲突
  49. self._cmd_queue = queue.Queue()
  50. self._cmd_min_interval = 0.10 # 相邻指令最小间隔 100ms
  51. self._last_cmd_time = 0
  52. self._cmd_worker_thread = threading.Thread(target=self._cmd_worker, daemon=True)
  53. self._cmd_worker_thread.start()
  54. # 启动后台重连监控线程,确保断线后可持续自动恢复
  55. self._start_reconnect_monitor()
  56. def _cmd_worker(self):
  57. """串口发送队列工作线程:统一按最小间隔逐个发送指令"""
  58. logger.info("启动串口发送队列工作线程")
  59. while True:
  60. try:
  61. item = self._cmd_queue.get()
  62. if item is None:
  63. break
  64. # 确保相邻指令最小间隔
  65. elapsed = time.time() - self._last_cmd_time
  66. if elapsed < self._cmd_min_interval:
  67. time.sleep(self._cmd_min_interval - elapsed)
  68. cmd_type = item.get('type')
  69. try:
  70. if cmd_type == 'send_and_wait':
  71. result = self._do_send_and_wait(
  72. item['data'],
  73. timeout=item.get('timeout', 2.0),
  74. min_response_bytes=item.get('min_response_bytes', 1)
  75. )
  76. elif cmd_type == 'wait_response':
  77. result = self._do_wait_response(
  78. timeout=item.get('timeout', 2.0),
  79. min_response_bytes=item.get('min_response_bytes', 1)
  80. )
  81. elif cmd_type == 'send_raw':
  82. result = self._do_send_raw(item['data'])
  83. elif cmd_type == 'send_data':
  84. result = self._do_send_data(item['data'], item.get('encoding', 'utf-8'))
  85. else:
  86. result = {'error': f'未知命令类型: {cmd_type}'}
  87. except Exception as e:
  88. logger.error(f"执行串口命令失败: {e}")
  89. result = {'error': str(e)}
  90. self._last_cmd_time = time.time()
  91. callback = item.get('callback')
  92. if callback:
  93. callback(result)
  94. except Exception as e:
  95. logger.error(f"串口发送队列工作线程异常: {e}")
  96. logger.info("串口发送队列工作线程结束")
  97. def _do_send_raw(self, data: bytes):
  98. """实际执行 send_raw(在队列工作线程中调用)"""
  99. try:
  100. with self.lock:
  101. if not self.is_connected or not self.ser or not self.ser.is_open:
  102. return False, "serial port not connected"
  103. bytes_sent = self.ser.write(data)
  104. self.ser.flush()
  105. if self.send_callback and data:
  106. self.send_callback(data.hex())
  107. return True, "send ok"
  108. except Exception as e:
  109. error_msg = f"send raw failed: {str(e)}"
  110. logger.error(error_msg)
  111. if self.error_callback:
  112. self.error_callback(error_msg)
  113. return False, error_msg
  114. def _do_send_data(self, data, encoding='utf-8'):
  115. """实际执行 send_data(在队列工作线程中调用)"""
  116. try:
  117. with self.lock:
  118. if not self.is_connected or not self.ser or not self.ser.is_open:
  119. return False, "串口未连接"
  120. # 确保数据以换行符结束
  121. if isinstance(data, str):
  122. if not data.endswith('\n'):
  123. data += '\n'
  124. bytes_data = data.encode(encoding)
  125. elif isinstance(data, bytes):
  126. if not data.endswith(b'\n'):
  127. bytes_data = data + b'\n'
  128. else:
  129. bytes_data = data
  130. else:
  131. raise TypeError("数据必须是字符串或字节类型")
  132. bytes_sent = self.ser.write(bytes_data)
  133. self.ser.flush()
  134. logger.debug(f"发送数据到串口: {bytes_data.hex()[:50]}... (共{bytes_sent}字节)")
  135. return True, "发送成功"
  136. except Exception as e:
  137. error_msg = f"发送失败: {str(e)}"
  138. logger.error(error_msg)
  139. if self.error_callback:
  140. self.error_callback(error_msg)
  141. return False, error_msg
  142. def _emit_received_data(self, data: bytes):
  143. """把接收到的原始数据推送给 data_callback(与后台读线程行为一致)"""
  144. if not data or not self.data_callback:
  145. return
  146. hex_data = data.hex()
  147. self.data_callback(hex_data)
  148. try:
  149. decoded = data.decode('utf-8').strip()
  150. if decoded:
  151. self.data_callback(decoded)
  152. except Exception:
  153. pass
  154. def _do_send_and_wait(self, data: bytes, timeout: float = 2.0, min_response_bytes: int = 1) -> bytes:
  155. """实际执行 send_and_wait(在队列工作线程中调用)"""
  156. # 进入同步接收模式,暂停后台读取线程
  157. with self._sync_lock:
  158. self._sync_receive_active = True
  159. try:
  160. # 阶段 1:发送数据并清空缓冲区
  161. with self.lock:
  162. if not self.ser or not self.ser.is_open:
  163. return b''
  164. self.raw_response_buffer.clear()
  165. self.ser.reset_input_buffer()
  166. self.ser.write(data)
  167. self.ser.flush()
  168. # 485 半双工:等最后字节发完并给收发器一点切换方向的时间
  169. time.sleep(0.03)
  170. if self.send_callback and data:
  171. self.send_callback(data.hex())
  172. # 阶段 2:等待响应。由于后台读线程已暂停,这里独占串口读取。
  173. response = b''
  174. start = time.time()
  175. while time.time() - start < timeout:
  176. # 先清掉缓冲区内可能残留的报文(正常情况下 sync 模式下读线程不会写入)
  177. with self.lock:
  178. while self.raw_response_buffer:
  179. try:
  180. response += bytes.fromhex(self.raw_response_buffer.pop(0))
  181. except Exception:
  182. pass
  183. try:
  184. if self.ser and self.ser.is_open and self.ser.in_waiting > 0:
  185. response += self.ser.read(self.ser.in_waiting)
  186. except Exception:
  187. pass
  188. if len(response) >= min_response_bytes:
  189. # 再等一小段时间收集可能的后续数据
  190. time.sleep(0.05)
  191. with self.lock:
  192. while self.raw_response_buffer:
  193. try:
  194. response += bytes.fromhex(self.raw_response_buffer.pop(0))
  195. except Exception:
  196. pass
  197. try:
  198. if self.ser and self.ser.is_open and self.ser.in_waiting > 0:
  199. response += self.ser.read(self.ser.in_waiting)
  200. except Exception:
  201. pass
  202. self._emit_received_data(response)
  203. return response
  204. time.sleep(0.01)
  205. self._emit_received_data(response)
  206. return response
  207. finally:
  208. # 退出同步接收模式,恢复后台读取线程
  209. with self._sync_lock:
  210. self._sync_receive_active = False
  211. def _do_wait_response(self, timeout: float = 2.0, min_response_bytes: int = 1) -> bytes:
  212. """实际执行 wait_response(在队列工作线程中调用)。
  213. 与 _do_send_and_wait 的区别:
  214. - 不发送任何数据(不 write),因此不触发 send_callback,不会在前端产生空 JSON;
  215. - 不清空接收缓冲(不 reset_input_buffer),避免把从机已回的应答冲掉。
  216. 用于发送完指令(send_raw)后专门等待从机应答。
  217. 收到 min_response_bytes 字节后,继续等到总线空闲(100ms 无新数据)再返回,
  218. 以便把回显与真实应答一并收齐后再交由调用方解析。
  219. """
  220. # 进入同步接收模式,暂停后台读取线程,独占串口读取
  221. with self._sync_lock:
  222. self._sync_receive_active = True
  223. try:
  224. response = b''
  225. start = time.time()
  226. last_data_time = start
  227. while time.time() - start < timeout:
  228. # 从后台读线程缓冲 + 串口硬件缓冲两处收集
  229. chunk = b''
  230. with self.lock:
  231. while self.raw_response_buffer:
  232. try:
  233. chunk += bytes.fromhex(self.raw_response_buffer.pop(0))
  234. except Exception:
  235. pass
  236. try:
  237. if self.ser and self.ser.is_open and self.ser.in_waiting > 0:
  238. chunk += self.ser.read(self.ser.in_waiting)
  239. except Exception:
  240. pass
  241. if chunk:
  242. response += chunk
  243. last_data_time = time.time()
  244. # 收够最小长度后,等总线空闲再返回,避免回显导致提前返回而漏掉真实应答
  245. if len(response) >= min_response_bytes:
  246. if time.time() - last_data_time > 0.1:
  247. self._emit_received_data(response)
  248. return response
  249. time.sleep(0.01)
  250. self._emit_received_data(response)
  251. return response
  252. finally:
  253. with self._sync_lock:
  254. self._sync_receive_active = False
  255. def list_ports(self):
  256. """列出系统中可用的串口"""
  257. ports = []
  258. try:
  259. # 首先尝试使用serial.tools.list_ports
  260. try:
  261. detected_ports = [port.device for port in serial.tools.list_ports.comports()]
  262. ports.extend(detected_ports)
  263. except Exception as e:
  264. logger.warning(f"使用serial.tools.list_ports失败: {str(e)}")
  265. # 根据不同平台进行补充查找
  266. system = platform.system()
  267. if system == 'Windows':
  268. try:
  269. import winreg
  270. # 在Windows系统中读取注册表
  271. key = winreg.OpenKey(winreg.HKEY_LOCAL_MACHINE,
  272. r'HARDWARE\DEVICEMAP\SERIALCOMM')
  273. i = 0
  274. while True:
  275. try:
  276. port, value, _ = winreg.EnumValue(key, i)
  277. if value not in ports:
  278. ports.append(value)
  279. i += 1
  280. except OSError:
  281. break
  282. except Exception as e:
  283. logger.error(f"读取Windows串口注册表失败: {str(e)}")
  284. elif system == 'Darwin': # macOS
  285. # 使用glob查找/dev/tty.*设备
  286. darwin_ports = glob.glob('/dev/tty.*')
  287. # 过滤掉不需要的端口
  288. for port in darwin_ports:
  289. if not ('Bluetooth' in port or 'debug' in port or 'com.apple' in port) and port not in ports:
  290. ports.append(port)
  291. elif system == 'Linux':
  292. # 使用glob查找Linux系统中的串口
  293. linux_ports = glob.glob('/dev/ttyS*') + glob.glob('/dev/ttyUSB*') + glob.glob('/dev/ttyACM*')
  294. for port in linux_ports:
  295. if port not in ports:
  296. ports.append(port)
  297. logger.info(f"找到 {len(ports)} 个可用串口: {ports}")
  298. except Exception as e:
  299. logger.error(f"列出串口时出错: {str(e)}")
  300. return sorted(ports) # 排序返回
  301. def connect(self, port, baudrate=9600, timeout=1, **kwargs):
  302. """连接到串口"""
  303. try:
  304. # 构建配置
  305. config = SerialConfig(
  306. port=port,
  307. baudrate=baudrate,
  308. timeout=timeout,
  309. **kwargs
  310. )
  311. with self.lock:
  312. if self.is_connected:
  313. # RLock 允许递归,避免死锁
  314. self.disconnect()
  315. # serial.Serial() 可能阻塞(端口被占用等),不要在持有锁的情况下调用,
  316. # 否则 disconnect() 会长时间等待锁而无法响应。
  317. logger.info(f"尝试连接串口: {port}, 波特率: {baudrate}")
  318. ser = serial.Serial(
  319. port=config.port,
  320. baudrate=config.baudrate,
  321. bytesize=config.bytesize,
  322. parity=config.parity,
  323. stopbits=config.stopbits,
  324. timeout=config.timeout,
  325. xonxoff=config.xonxoff,
  326. rtscts=config.rtscts,
  327. dsrdtr=config.dsrdtr
  328. )
  329. # 检查连接是否成功
  330. if not ser.is_open:
  331. ser.close()
  332. raise Exception("串口打开失败")
  333. with self.lock:
  334. self.ser = ser
  335. self.is_connected = True
  336. self.stop_event.clear()
  337. self.current_config = config
  338. self.reconnect_attempts = 0
  339. # 启动读取线程
  340. self.read_thread = threading.Thread(target=self._read_loop, daemon=True)
  341. self.read_thread.start()
  342. if self.status_callback:
  343. self.status_callback(True)
  344. logger.info(f"已连接到 {port},波特率 {baudrate}")
  345. return True, f"已连接到 {port},波特率 {baudrate}"
  346. except Exception as e:
  347. error_msg = f"连接失败: {str(e)}"
  348. logger.error(error_msg)
  349. if self.status_callback:
  350. self.status_callback(False)
  351. if self.error_callback:
  352. self.error_callback(error_msg)
  353. return False, error_msg
  354. def disconnect(self):
  355. """断开串口连接"""
  356. try:
  357. logger.info("断开串口连接")
  358. # 阶段 0:清空发送队列中待执行的指令,避免断开后再发数据
  359. try:
  360. while not self._cmd_queue.empty():
  361. self._cmd_queue.get_nowait()
  362. except Exception:
  363. pass
  364. # 阶段 1:通知读取线程退出,不要在持有锁的情况下 join,
  365. # 否则读取线程异常时调用 _close_on_error 会拿不到锁而死锁。
  366. with self.lock:
  367. self.stop_event.set()
  368. if self.read_thread and self.read_thread.is_alive():
  369. self.read_thread.join(timeout=2.0)
  370. if self.read_thread.is_alive():
  371. logger.warning("读取线程未能正常终止")
  372. # 阶段 2:关闭串口并清理状态
  373. with self.lock:
  374. if self.ser and self.ser.is_open:
  375. try:
  376. self.ser.close()
  377. except Exception as e:
  378. logger.error(f"关闭串口时出错: {str(e)}")
  379. self.ser = None
  380. self.is_connected = False
  381. self.current_config = None
  382. if self.status_callback:
  383. self.status_callback(False)
  384. return True, "已断开连接"
  385. except Exception as e:
  386. error_msg = f"断开连接失败: {str(e)}"
  387. logger.error(error_msg)
  388. if self.error_callback:
  389. self.error_callback(error_msg)
  390. return False, error_msg
  391. def _read_loop(self):
  392. """读取串口数据的循环"""
  393. logger.info("启动串口读取线程")
  394. while not self.stop_event.is_set():
  395. try:
  396. # send_and_wait 正在同步接收时,后台读线程让出总线,
  397. # 避免读线程把响应抢走导致同步调用超时。
  398. with self._sync_lock:
  399. sync_active = self._sync_receive_active
  400. if sync_active:
  401. time.sleep(0.001)
  402. continue
  403. if self.ser and self.ser.is_open:
  404. # 使用in_waiting提高效率
  405. if self.ser.in_waiting > 0:
  406. data = self.ser.read(self.ser.in_waiting)
  407. # 默认将原始数据以十六进制存入缓冲区
  408. hex_data = data.hex()
  409. with self.lock:
  410. self.raw_response_buffer.append(hex_data)
  411. if self.data_callback:
  412. self.data_callback(hex_data)
  413. # 如果能解码为文本,也通知回调
  414. try:
  415. decoded_data = data.decode('utf-8').strip()
  416. if decoded_data and self.data_callback:
  417. self.data_callback(decoded_data)
  418. except:
  419. pass
  420. time.sleep(0.001)
  421. except Exception as e:
  422. error_msg = f"读取串口数据错误: {str(e)}"
  423. logger.error(error_msg)
  424. if self.error_callback:
  425. self.error_callback(error_msg)
  426. # 读取线程只负责关闭当前连接并通知,重连由独立监控线程负责,
  427. # 避免在读取线程内调用 connect() 造成自连接/自 join 的问题
  428. self._close_on_error()
  429. break
  430. # 线程结束时清理资源
  431. logger.info("串口读取线程结束")
  432. def _close_on_error(self):
  433. """读取异常时关闭串口并触发状态回调,但不直接重连"""
  434. with self.lock:
  435. if self.ser and self.ser.is_open:
  436. try:
  437. self.ser.close()
  438. except Exception:
  439. pass
  440. self.ser = None
  441. self.is_connected = False
  442. if self.status_callback:
  443. try:
  444. self.status_callback(False)
  445. except Exception:
  446. pass
  447. logger.warning("串口因读取错误已关闭,等待重连监控线程恢复")
  448. def _start_reconnect_monitor(self):
  449. """启动独立后台线程,在串口断开时持续尝试重连"""
  450. if getattr(self, '_reconnect_monitor_thread', None) and self._reconnect_monitor_thread.is_alive():
  451. return
  452. self._reconnect_monitor_thread = threading.Thread(target=self._reconnect_monitor, daemon=True)
  453. self._reconnect_monitor_thread.start()
  454. def _reconnect_monitor(self):
  455. """后台重连监控:只要保存过配置就无限重试,成功则重置计数"""
  456. logger.info("启动串口重连监控线程")
  457. while True:
  458. try:
  459. with self.lock:
  460. connected = self.is_connected
  461. config = self.current_config
  462. if not connected and config is not None:
  463. self.reconnect_attempts += 1
  464. attempt = self.reconnect_attempts
  465. logger.warning(f"重连监控尝试连接串口... (第{attempt}次)")
  466. success, msg = self.connect(
  467. port=config.port,
  468. baudrate=config.baudrate,
  469. timeout=config.timeout,
  470. bytesize=config.bytesize,
  471. parity=config.parity,
  472. stopbits=config.stopbits,
  473. xonxoff=config.xonxoff,
  474. rtscts=config.rtscts,
  475. dsrdtr=config.dsrdtr
  476. )
  477. if not success:
  478. # 指数退避,最长 30 秒
  479. delay = min(30, 2 + attempt * 2)
  480. logger.warning(f"重连失败: {msg},{delay}秒后再次尝试")
  481. time.sleep(delay)
  482. else:
  483. logger.info("串口重连成功")
  484. else:
  485. # 已连接或未保存配置时,重置失败计数并降低检查频率
  486. if connected:
  487. self.reconnect_attempts = 0
  488. time.sleep(3)
  489. except Exception as e:
  490. logger.error(f"重连监控线程异常: {e}")
  491. time.sleep(5)
  492. def send_data(self, data, encoding='utf-8'):
  493. """发送数据到串口(入队,由队列工作线程统一发送)"""
  494. result_container = {}
  495. event = threading.Event()
  496. def callback(result):
  497. result_container['result'] = result
  498. event.set()
  499. self._cmd_queue.put({
  500. 'type': 'send_data',
  501. 'data': data,
  502. 'encoding': encoding,
  503. 'callback': callback
  504. })
  505. # 最多等待 5 秒,避免队列卡死导致调用方永远阻塞
  506. if not event.wait(5):
  507. return False, "串口发送队列超时"
  508. return result_container.get('result', (False, "未知错误"))
  509. def set_data_callback(self, callback):
  510. """设置数据接收回调函数"""
  511. self.data_callback = callback
  512. def set_send_callback(self, callback):
  513. """设置数据发送回调函数"""
  514. self.send_callback = callback
  515. def set_status_callback(self, callback):
  516. """设置状态变化回调函数"""
  517. self.status_callback = callback
  518. def set_error_callback(self, callback):
  519. """设置错误回调函数"""
  520. self.error_callback = callback
  521. def get_status(self):
  522. """获取当前连接状态"""
  523. with self.lock:
  524. return {
  525. 'connected': self.is_connected,
  526. 'config': self.current_config,
  527. 'has_error': self.reconnect_attempts > 0
  528. }
  529. def _should_reconnect(self):
  530. """判断是否应该尝试重连"""
  531. self.reconnect_attempts += 1
  532. return self.reconnect_attempts <= self.max_reconnect_attempts
  533. def send_raw(self, data: bytes):
  534. """send raw binary data without adding newline(入队,由队列工作线程统一发送)"""
  535. result_container = {}
  536. event = threading.Event()
  537. def callback(result):
  538. result_container['result'] = result
  539. event.set()
  540. self._cmd_queue.put({
  541. 'type': 'send_raw',
  542. 'data': data,
  543. 'callback': callback
  544. })
  545. # 最多等待 5 秒,避免队列卡死导致调用方永远阻塞
  546. if not event.wait(5):
  547. return False, "串口发送队列超时"
  548. return result_container.get('result', (False, "未知错误"))
  549. def flush_input(self):
  550. """清空输入缓冲区"""
  551. try:
  552. with self.lock:
  553. if self.ser and self.ser.is_open:
  554. self.ser.reset_input_buffer()
  555. return True, "输入缓冲区已清空"
  556. return False, "串口未连接"
  557. except Exception as e:
  558. error_msg = f"清空缓冲区失败: {str(e)}"
  559. logger.error(error_msg)
  560. return False, error_msg
  561. def flush_output(self):
  562. """清空输出缓冲区"""
  563. try:
  564. with self.lock:
  565. if self.ser and self.ser.is_open:
  566. self.ser.reset_output_buffer()
  567. return True, "输出缓冲区已清空"
  568. return False, "串口未连接"
  569. except Exception as e:
  570. error_msg = f"清空缓冲区失败: {str(e)}"
  571. logger.error(error_msg)
  572. return False, error_msg
  573. def send_and_wait(self, data: bytes, timeout: float = 2.0, min_response_bytes: int = 1) -> bytes:
  574. """sync send and wait(入队,由队列工作线程统一发送并等待响应)"""
  575. result_container = {}
  576. event = threading.Event()
  577. def callback(result):
  578. result_container['result'] = result
  579. event.set()
  580. self._cmd_queue.put({
  581. 'type': 'send_and_wait',
  582. 'data': data,
  583. 'timeout': timeout,
  584. 'min_response_bytes': min_response_bytes,
  585. 'callback': callback
  586. })
  587. # 最多等待 timeout + 队列处理余量
  588. wait_time = timeout + 5
  589. if not event.wait(wait_time):
  590. return b''
  591. return result_container.get('result', b'')
  592. def wait_response(self, timeout: float = 2.0, min_response_bytes: int = 1) -> bytes:
  593. """只读取从机应答,不发送任何数据(入队,由队列工作线程统一执行)。
  594. 与 send_and_wait 的区别:不 write 数据、不触发 send_callback、不清空接收缓冲。
  595. 用于发送完指令(send_raw)后等待从机应答,避免发空数据污染前端显示、
  596. 避免清缓冲冲掉已到达的应答。
  597. """
  598. result_container = {}
  599. event = threading.Event()
  600. def callback(result):
  601. result_container['result'] = result
  602. event.set()
  603. self._cmd_queue.put({
  604. 'type': 'wait_response',
  605. 'timeout': timeout,
  606. 'min_response_bytes': min_response_bytes,
  607. 'callback': callback
  608. })
  609. # 最多等待 timeout + 队列处理余量
  610. wait_time = timeout + 5
  611. if not event.wait(wait_time):
  612. return b''
  613. return result_container.get('result', b'')