| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620 |
- import serial
- import serial.tools.list_ports
- import threading
- import time
- import logging
- import platform
- import glob
- import queue
- from dataclasses import dataclass
- # 配置日志
- logging.basicConfig(
- level=logging.INFO,
- format='%(asctime)s - %(name)s - %(levelname)s - %(message)s'
- )
- logger = logging.getLogger('serial_port')
- @dataclass
- class SerialConfig:
- """串口配置数据类"""
- port: str
- baudrate: int = 115200
- bytesize: int = serial.EIGHTBITS
- parity: str = serial.PARITY_NONE
- stopbits: int = serial.STOPBITS_ONE
- timeout: float = 1.0
- xonxoff: bool = False
- rtscts: bool = False
- dsrdtr: bool = False
- class SerialPort:
- """串口通信类,提供串口连接、读写和状态管理功能"""
-
- def __init__(self):
- self.ser = None
- self.is_connected = False
- self.lock = threading.RLock() # 使用可重入锁
- self.read_thread = None
- self.stop_event = threading.Event()
- self.data_callback = None
- self.send_callback = None
- self.status_callback = None
- self.error_callback = None
- self.current_config = None
- self.reconnect_attempts = 0
- self.max_reconnect_attempts = 3
- self.raw_response_buffer = []
- self._reconnect_monitor_thread = None
- # 同步接收模式:send_and_wait 活动时暂停后台读取线程,防止响应被读线程抢走
- self._sync_receive_active = False
- self._sync_lock = threading.Lock()
- # 串口发送队列:所有发送指令统一排队,按最小间隔逐个发送,避免总线冲突
- self._cmd_queue = queue.Queue()
- self._cmd_min_interval = 0.10 # 相邻指令最小间隔 100ms
- self._last_cmd_time = 0
- self._cmd_worker_thread = threading.Thread(target=self._cmd_worker, daemon=True)
- self._cmd_worker_thread.start()
- # 启动后台重连监控线程,确保断线后可持续自动恢复
- self._start_reconnect_monitor()
- def _cmd_worker(self):
- """串口发送队列工作线程:统一按最小间隔逐个发送指令"""
- logger.info("启动串口发送队列工作线程")
- while True:
- try:
- item = self._cmd_queue.get()
- if item is None:
- break
- # 确保相邻指令最小间隔
- elapsed = time.time() - self._last_cmd_time
- if elapsed < self._cmd_min_interval:
- time.sleep(self._cmd_min_interval - elapsed)
- cmd_type = item.get('type')
- try:
- if cmd_type == 'send_and_wait':
- result = self._do_send_and_wait(
- item['data'],
- timeout=item.get('timeout', 2.0),
- min_response_bytes=item.get('min_response_bytes', 1)
- )
- elif cmd_type == 'send_raw':
- result = self._do_send_raw(item['data'])
- elif cmd_type == 'send_data':
- result = self._do_send_data(item['data'], item.get('encoding', 'utf-8'))
- else:
- result = {'error': f'未知命令类型: {cmd_type}'}
- except Exception as e:
- logger.error(f"执行串口命令失败: {e}")
- result = {'error': str(e)}
- self._last_cmd_time = time.time()
- callback = item.get('callback')
- if callback:
- callback(result)
- except Exception as e:
- logger.error(f"串口发送队列工作线程异常: {e}")
- logger.info("串口发送队列工作线程结束")
- def _do_send_raw(self, data: bytes):
- """实际执行 send_raw(在队列工作线程中调用)"""
- try:
- with self.lock:
- if not self.is_connected or not self.ser or not self.ser.is_open:
- return False, "serial port not connected"
- bytes_sent = self.ser.write(data)
- self.ser.flush()
- if self.send_callback:
- self.send_callback(data.hex())
- return True, "send ok"
- except Exception as e:
- error_msg = f"send raw failed: {str(e)}"
- logger.error(error_msg)
- if self.error_callback:
- self.error_callback(error_msg)
- return False, error_msg
- def _do_send_data(self, data, encoding='utf-8'):
- """实际执行 send_data(在队列工作线程中调用)"""
- try:
- with self.lock:
- if not self.is_connected or not self.ser or not self.ser.is_open:
- return False, "串口未连接"
- # 确保数据以换行符结束
- if isinstance(data, str):
- if not data.endswith('\n'):
- data += '\n'
- bytes_data = data.encode(encoding)
- elif isinstance(data, bytes):
- if not data.endswith(b'\n'):
- bytes_data = data + b'\n'
- else:
- bytes_data = data
- else:
- raise TypeError("数据必须是字符串或字节类型")
- bytes_sent = self.ser.write(bytes_data)
- self.ser.flush()
- logger.debug(f"发送数据到串口: {bytes_data.hex()[:50]}... (共{bytes_sent}字节)")
- return True, "发送成功"
- except Exception as e:
- error_msg = f"发送失败: {str(e)}"
- logger.error(error_msg)
- if self.error_callback:
- self.error_callback(error_msg)
- return False, error_msg
- def _emit_received_data(self, data: bytes):
- """把接收到的原始数据推送给 data_callback(与后台读线程行为一致)"""
- if not data or not self.data_callback:
- return
- hex_data = data.hex()
- self.data_callback(hex_data)
- try:
- decoded = data.decode('utf-8').strip()
- if decoded:
- self.data_callback(decoded)
- except Exception:
- pass
- def _do_send_and_wait(self, data: bytes, timeout: float = 2.0, min_response_bytes: int = 1) -> bytes:
- """实际执行 send_and_wait(在队列工作线程中调用)"""
- # 进入同步接收模式,暂停后台读取线程
- with self._sync_lock:
- self._sync_receive_active = True
- try:
- # 阶段 1:发送数据并清空缓冲区
- with self.lock:
- if not self.ser or not self.ser.is_open:
- return b''
- self.raw_response_buffer.clear()
- self.ser.reset_input_buffer()
- self.ser.write(data)
- self.ser.flush()
- # 485 半双工:等最后字节发完并给收发器一点切换方向的时间
- time.sleep(0.03)
- if self.send_callback:
- self.send_callback(data.hex())
- # 阶段 2:等待响应。由于后台读线程已暂停,这里独占串口读取。
- response = b''
- start = time.time()
- while time.time() - start < timeout:
- # 先清掉缓冲区内可能残留的报文(正常情况下 sync 模式下读线程不会写入)
- with self.lock:
- while self.raw_response_buffer:
- try:
- response += bytes.fromhex(self.raw_response_buffer.pop(0))
- except Exception:
- pass
- try:
- if self.ser and self.ser.is_open and self.ser.in_waiting > 0:
- response += self.ser.read(self.ser.in_waiting)
- except Exception:
- pass
- if len(response) >= min_response_bytes:
- # 再等一小段时间收集可能的后续数据
- time.sleep(0.05)
- with self.lock:
- while self.raw_response_buffer:
- try:
- response += bytes.fromhex(self.raw_response_buffer.pop(0))
- except Exception:
- pass
- try:
- if self.ser and self.ser.is_open and self.ser.in_waiting > 0:
- response += self.ser.read(self.ser.in_waiting)
- except Exception:
- pass
- self._emit_received_data(response)
- return response
- time.sleep(0.01)
- self._emit_received_data(response)
- return response
- finally:
- # 退出同步接收模式,恢复后台读取线程
- with self._sync_lock:
- self._sync_receive_active = False
- def list_ports(self):
- """列出系统中可用的串口"""
- ports = []
-
- try:
- # 首先尝试使用serial.tools.list_ports
- try:
- detected_ports = [port.device for port in serial.tools.list_ports.comports()]
- ports.extend(detected_ports)
- except Exception as e:
- logger.warning(f"使用serial.tools.list_ports失败: {str(e)}")
-
- # 根据不同平台进行补充查找
- system = platform.system()
-
- if system == 'Windows':
- try:
- import winreg
- # 在Windows系统中读取注册表
- key = winreg.OpenKey(winreg.HKEY_LOCAL_MACHINE,
- r'HARDWARE\DEVICEMAP\SERIALCOMM')
- i = 0
- while True:
- try:
- port, value, _ = winreg.EnumValue(key, i)
- if value not in ports:
- ports.append(value)
- i += 1
- except OSError:
- break
- except Exception as e:
- logger.error(f"读取Windows串口注册表失败: {str(e)}")
-
- elif system == 'Darwin': # macOS
- # 使用glob查找/dev/tty.*设备
- darwin_ports = glob.glob('/dev/tty.*')
- # 过滤掉不需要的端口
- for port in darwin_ports:
- if not ('Bluetooth' in port or 'debug' in port or 'com.apple' in port) and port not in ports:
- ports.append(port)
-
- elif system == 'Linux':
- # 使用glob查找Linux系统中的串口
- linux_ports = glob.glob('/dev/ttyS*') + glob.glob('/dev/ttyUSB*') + glob.glob('/dev/ttyACM*')
- for port in linux_ports:
- if port not in ports:
- ports.append(port)
-
- logger.info(f"找到 {len(ports)} 个可用串口: {ports}")
- except Exception as e:
- logger.error(f"列出串口时出错: {str(e)}")
-
- return sorted(ports) # 排序返回
-
- def connect(self, port, baudrate=9600, timeout=1, **kwargs):
- """连接到串口"""
- try:
- # 构建配置
- config = SerialConfig(
- port=port,
- baudrate=baudrate,
- timeout=timeout,
- **kwargs
- )
- with self.lock:
- if self.is_connected:
- # RLock 允许递归,避免死锁
- self.disconnect()
- # serial.Serial() 可能阻塞(端口被占用等),不要在持有锁的情况下调用,
- # 否则 disconnect() 会长时间等待锁而无法响应。
- logger.info(f"尝试连接串口: {port}, 波特率: {baudrate}")
- ser = serial.Serial(
- port=config.port,
- baudrate=config.baudrate,
- bytesize=config.bytesize,
- parity=config.parity,
- stopbits=config.stopbits,
- timeout=config.timeout,
- xonxoff=config.xonxoff,
- rtscts=config.rtscts,
- dsrdtr=config.dsrdtr
- )
- # 检查连接是否成功
- if not ser.is_open:
- ser.close()
- raise Exception("串口打开失败")
- with self.lock:
- self.ser = ser
- self.is_connected = True
- self.stop_event.clear()
- self.current_config = config
- self.reconnect_attempts = 0
- # 启动读取线程
- self.read_thread = threading.Thread(target=self._read_loop, daemon=True)
- self.read_thread.start()
- if self.status_callback:
- self.status_callback(True)
- logger.info(f"已连接到 {port},波特率 {baudrate}")
- return True, f"已连接到 {port},波特率 {baudrate}"
- except Exception as e:
- error_msg = f"连接失败: {str(e)}"
- logger.error(error_msg)
- if self.status_callback:
- self.status_callback(False)
- if self.error_callback:
- self.error_callback(error_msg)
- return False, error_msg
-
- def disconnect(self):
- """断开串口连接"""
- try:
- logger.info("断开串口连接")
- # 阶段 0:清空发送队列中待执行的指令,避免断开后再发数据
- try:
- while not self._cmd_queue.empty():
- self._cmd_queue.get_nowait()
- except Exception:
- pass
- # 阶段 1:通知读取线程退出,不要在持有锁的情况下 join,
- # 否则读取线程异常时调用 _close_on_error 会拿不到锁而死锁。
- with self.lock:
- self.stop_event.set()
- if self.read_thread and self.read_thread.is_alive():
- self.read_thread.join(timeout=2.0)
- if self.read_thread.is_alive():
- logger.warning("读取线程未能正常终止")
- # 阶段 2:关闭串口并清理状态
- with self.lock:
- if self.ser and self.ser.is_open:
- try:
- self.ser.close()
- except Exception as e:
- logger.error(f"关闭串口时出错: {str(e)}")
- self.ser = None
- self.is_connected = False
- self.current_config = None
- if self.status_callback:
- self.status_callback(False)
- return True, "已断开连接"
- except Exception as e:
- error_msg = f"断开连接失败: {str(e)}"
- logger.error(error_msg)
- if self.error_callback:
- self.error_callback(error_msg)
- return False, error_msg
-
- def _read_loop(self):
- """读取串口数据的循环"""
- logger.info("启动串口读取线程")
- while not self.stop_event.is_set():
- try:
- # send_and_wait 正在同步接收时,后台读线程让出总线,
- # 避免读线程把响应抢走导致同步调用超时。
- with self._sync_lock:
- sync_active = self._sync_receive_active
- if sync_active:
- time.sleep(0.001)
- continue
- if self.ser and self.ser.is_open:
- # 使用in_waiting提高效率
- if self.ser.in_waiting > 0:
- data = self.ser.read(self.ser.in_waiting)
- # 默认将原始数据以十六进制存入缓冲区
- hex_data = data.hex()
- with self.lock:
- self.raw_response_buffer.append(hex_data)
- if self.data_callback:
- self.data_callback(hex_data)
- # 如果能解码为文本,也通知回调
- try:
- decoded_data = data.decode('utf-8').strip()
- if decoded_data and self.data_callback:
- self.data_callback(decoded_data)
- except:
- pass
- time.sleep(0.001)
- except Exception as e:
- error_msg = f"读取串口数据错误: {str(e)}"
- logger.error(error_msg)
- if self.error_callback:
- self.error_callback(error_msg)
- # 读取线程只负责关闭当前连接并通知,重连由独立监控线程负责,
- # 避免在读取线程内调用 connect() 造成自连接/自 join 的问题
- self._close_on_error()
- break
- # 线程结束时清理资源
- logger.info("串口读取线程结束")
- def _close_on_error(self):
- """读取异常时关闭串口并触发状态回调,但不直接重连"""
- with self.lock:
- if self.ser and self.ser.is_open:
- try:
- self.ser.close()
- except Exception:
- pass
- self.ser = None
- self.is_connected = False
- if self.status_callback:
- try:
- self.status_callback(False)
- except Exception:
- pass
- logger.warning("串口因读取错误已关闭,等待重连监控线程恢复")
- def _start_reconnect_monitor(self):
- """启动独立后台线程,在串口断开时持续尝试重连"""
- if getattr(self, '_reconnect_monitor_thread', None) and self._reconnect_monitor_thread.is_alive():
- return
- self._reconnect_monitor_thread = threading.Thread(target=self._reconnect_monitor, daemon=True)
- self._reconnect_monitor_thread.start()
- def _reconnect_monitor(self):
- """后台重连监控:只要保存过配置就无限重试,成功则重置计数"""
- logger.info("启动串口重连监控线程")
- while True:
- try:
- with self.lock:
- connected = self.is_connected
- config = self.current_config
- if not connected and config is not None:
- self.reconnect_attempts += 1
- attempt = self.reconnect_attempts
- logger.warning(f"重连监控尝试连接串口... (第{attempt}次)")
- success, msg = self.connect(
- port=config.port,
- baudrate=config.baudrate,
- timeout=config.timeout,
- bytesize=config.bytesize,
- parity=config.parity,
- stopbits=config.stopbits,
- xonxoff=config.xonxoff,
- rtscts=config.rtscts,
- dsrdtr=config.dsrdtr
- )
- if not success:
- # 指数退避,最长 30 秒
- delay = min(30, 2 + attempt * 2)
- logger.warning(f"重连失败: {msg},{delay}秒后再次尝试")
- time.sleep(delay)
- else:
- logger.info("串口重连成功")
- else:
- # 已连接或未保存配置时,重置失败计数并降低检查频率
- if connected:
- self.reconnect_attempts = 0
- time.sleep(3)
- except Exception as e:
- logger.error(f"重连监控线程异常: {e}")
- time.sleep(5)
-
- def send_data(self, data, encoding='utf-8'):
- """发送数据到串口(入队,由队列工作线程统一发送)"""
- result_container = {}
- event = threading.Event()
- def callback(result):
- result_container['result'] = result
- event.set()
- self._cmd_queue.put({
- 'type': 'send_data',
- 'data': data,
- 'encoding': encoding,
- 'callback': callback
- })
- # 最多等待 5 秒,避免队列卡死导致调用方永远阻塞
- if not event.wait(5):
- return False, "串口发送队列超时"
- return result_container.get('result', (False, "未知错误"))
-
- def set_data_callback(self, callback):
- """设置数据接收回调函数"""
- self.data_callback = callback
- def set_send_callback(self, callback):
- """设置数据发送回调函数"""
- self.send_callback = callback
-
- def set_status_callback(self, callback):
- """设置状态变化回调函数"""
- self.status_callback = callback
-
- def set_error_callback(self, callback):
- """设置错误回调函数"""
- self.error_callback = callback
-
- def get_status(self):
- """获取当前连接状态"""
- with self.lock:
- return {
- 'connected': self.is_connected,
- 'config': self.current_config,
- 'has_error': self.reconnect_attempts > 0
- }
-
- def _should_reconnect(self):
- """判断是否应该尝试重连"""
- self.reconnect_attempts += 1
- return self.reconnect_attempts <= self.max_reconnect_attempts
-
- def send_raw(self, data: bytes):
- """send raw binary data without adding newline(入队,由队列工作线程统一发送)"""
- result_container = {}
- event = threading.Event()
- def callback(result):
- result_container['result'] = result
- event.set()
- self._cmd_queue.put({
- 'type': 'send_raw',
- 'data': data,
- 'callback': callback
- })
- # 最多等待 5 秒,避免队列卡死导致调用方永远阻塞
- if not event.wait(5):
- return False, "串口发送队列超时"
- return result_container.get('result', (False, "未知错误"))
- def flush_input(self):
- """清空输入缓冲区"""
- try:
- with self.lock:
- if self.ser and self.ser.is_open:
- self.ser.reset_input_buffer()
- return True, "输入缓冲区已清空"
- return False, "串口未连接"
- except Exception as e:
- error_msg = f"清空缓冲区失败: {str(e)}"
- logger.error(error_msg)
- return False, error_msg
-
- def flush_output(self):
- """清空输出缓冲区"""
- try:
- with self.lock:
- if self.ser and self.ser.is_open:
- self.ser.reset_output_buffer()
- return True, "输出缓冲区已清空"
- return False, "串口未连接"
- except Exception as e:
- error_msg = f"清空缓冲区失败: {str(e)}"
- logger.error(error_msg)
- return False, error_msg
- def send_and_wait(self, data: bytes, timeout: float = 2.0, min_response_bytes: int = 1) -> bytes:
- """sync send and wait(入队,由队列工作线程统一发送并等待响应)"""
- result_container = {}
- event = threading.Event()
- def callback(result):
- result_container['result'] = result
- event.set()
- self._cmd_queue.put({
- 'type': 'send_and_wait',
- 'data': data,
- 'timeout': timeout,
- 'min_response_bytes': min_response_bytes,
- 'callback': callback
- })
- # 最多等待 timeout + 队列处理余量
- wait_time = timeout + 5
- if not event.wait(wait_time):
- return b''
- return result_container.get('result', b'')
|