import datetime import queue import socket import threading import time from typing import List, Tuple, Optional, Dict, Any import codecs from bx import BxResp, NewBxDataPackCmd, NewBxAreaProgram, BxCmdFactory, NewBxFile # 核心常量定义 class Color: DEFAULT = 0 RED = 1 GREEN = 2 YELLOW = 3 BLUE = 4 LIGHT_BLUE = 5 LIGHT_PURPLE = 6 WHITE = 7 class RunMode: LOOP = 0 # 循环 LOOP_AND_STAY_AT_END = 1 # 循环直到最后,停留在最后一个动态区 LOOP_AND_TIMEOUT_OFF = 2 # 循环直到超时,超时后未更新不在显示 LOOP_AND_STAY_AT_LOGO = 3 # 循环完后,停留显示LOGO LOOP_AND_OFF = 4 # 循环完后不在显示 LOOP_AND_COUNT_OFF = 5 # 循环设定次数后不在显示 DEFAULT = LOOP # 默认 class DisplayMode: STATIC = 1 # 静止显示 QUICK_PUNCH = 2 # 快速打出 MOVE_LEFT = 3 # 向左移动 MOVE_RIGHT = 4 # 向右移动 MOVE_UP = 5 # 向上移动 MOVE_DOWN = 6 # 向下移动 FLICKER = 7 # 闪烁 DEFAULT = STATIC # 默认 class FlashFile: def __init__(self): self.msg = "" self.sound_data = "" self.color = Color.DEFAULT self.run_mode = RunMode.DEFAULT self.disp_mode = DisplayMode.DEFAULT self.origin_x = 0 self.origin_y = 0 self.width = 0 self.height = 0 def set_msg(self, msg: str, color: int): self.msg = msg self.color = color def set_mode(self, run_mode: int, disp_mode: int): self.run_mode = run_mode self.disp_mode = disp_mode def set_origin(self, x: int, x_is_pixel: bool, y: int): self.origin_y = y if x_is_pixel: self.origin_x = 0x8000 | x else: self.origin_x = x def set_area(self, w: int, y_is_pixel: bool, h: int): self.height = h if y_is_pixel: self.width = 0x8000 | w else: self.width = w class CommandQueue: def __init__(self, max_size: int = 100): self.queue = queue.Queue(maxsize=max_size) self.stop_flag = False self.thread: Optional[threading.Thread] = None def start(self, handle_func): """启动队列消费线程""" self.stop_flag = False self.thread = threading.Thread(target=self._consumer, args=(handle_func,), daemon=True) self.thread.start() def stop(self): """停止队列消费""" self.stop_flag = True if self.thread: self.thread.join(timeout=2) def add_task(self, data: bytes): """添加发送任务到队列""" try: self.queue.put(data, block=False) except queue.Full: print("发送队列已满,丢弃任务") def _consumer(self, handle_func): """队列消费逻辑""" while not self.stop_flag: try: data = self.queue.get(timeout=1) handle_func(data) self.queue.task_done() except queue.Empty: continue class Screen: def __init__(self, name: str, ip: str, port: str): self.name = name self.addr = f"{ip}:{port}" self.conn: Optional[socket.socket] = None self.live_state = False self.state_info = dict() # 替代bx.StateInfo self.params = dict() # 替代bx.Params self.send_queue = CommandQueue(max_size=100) # 新增:重连锁,防止多线程同时重连 self.reconnect_lock = threading.Lock() # 新增:退出标记,用于停止监控线程 self.stop_flag = False # 启动发送队列 self.send_queue.start(self.handle_send) # 启动连接监控线程 self.monitor_thread = threading.Thread(target=self._monitor_connection, daemon=True) self.monitor_thread.start() # 初始重连(非阻塞式,在后台尝试连接) self._initial_connect() def _initial_connect(self): """初始连接(非阻塞,快速尝试一次)""" try: ip, port = self.addr.split(":") conn = socket.create_connection((ip, int(port)), timeout=10) # 新增:设置Socket超时,避免recv无限阻塞 conn.settimeout(10.0) print(f"[{self.name}] 初始连接成功: {self.addr}") self.set_conn(conn) except Exception as e: print(f"[{self.name}] 初始连接失败: {self.addr}, 将在后台重试...") # 连接失败不阻塞,监控线程会自动重连 def handle_send(self, data: bytes): """实际执行TCP发送(由队列线程调用)""" if not self.get_live_state() or self.conn is None: return try: self.conn.sendall(data) except Exception as e: print(f"[{self.name}] TCP发送失败: {e}") # 优化:仅在发送失败时标记为断开,加锁保护 with self.reconnect_lock: self.live_state = False # ========== 新增:直接发送方法(跳过队列,用于雷达屏幕) ========== def send_direct(self, data: bytes): """ 直接发送数据(跳过队列) 用于强实时性场景(如雷达测速) """ if not self.get_live_state() or self.conn is None: return False try: self.conn.sendall(data) return True except Exception as e: print(f"[{self.name}] 直接发送失败: {e}") # 发送失败时标记为断开 with self.reconnect_lock: self.live_state = False return False def send(self, data: bytes): """添加发送任务到队列""" if not self.get_live_state(): return self.send_queue.add_task(data) def close(self): """关闭连接和队列""" # 新增:设置退出标记,停止监控线程 self.stop_flag = True self.send_queue.stop() # 加锁保护,防止并发修改 with self.reconnect_lock: if self.conn: try: self.conn.close() except Exception as e: print(f"[{self.name}] 关闭连接失败: {e}") self.conn = None self.live_state = False # 等待监控线程退出 if self.monitor_thread: self.monitor_thread.join(timeout=3) def get_live_state(self) -> bool: """获取在线状态(加锁保护)""" with self.reconnect_lock: return self.live_state def set_conn(self, conn: socket.socket): """设置连接并更新在线状态(加锁保护)""" with self.reconnect_lock: self.conn = conn # 新增:设置Socket超时 self.conn.settimeout(10.0) self.live_state = True def reconnect(self): """重连屏幕TCP连接(加锁防止并发)""" # 加锁:同一时间只能有一个重连操作 with self.reconnect_lock: if self.live_state: return # 关闭旧连接 if self.conn: try: self.conn.close() except Exception: pass self.conn = None # 新建连接 try: ip, port = self.addr.split(":") conn = socket.create_connection((ip, int(port)), timeout=10) conn.settimeout(10.0) print(f"[{self.name}] 连接成功: {self.addr}") self.set_conn(conn) # 优化:延迟读取状态/参数,给屏幕缓冲时间 time.sleep(0.5) # 读取屏状态和参数(非阻塞,失败不影响连接状态) try: state_resp = self.state() if state_resp: self.state_info = self._parse_state(state_resp.Data) params_resp = self.param() if params_resp: self.params = self._parse_params(params_resp.Data) except Exception as e: print(f"[{self.name}] 读取状态/参数失败(不影响连接): {e}") except Exception as e: print(f"[{self.name}] 重连失败: {self.addr}, error: {e}") self.live_state = False def force_reconnect(self, max_retries: int = 10): """强制重连屏幕TCP连接(带最大重试次数限制)""" retry_count = 0 while not self.get_live_state() and retry_count < max_retries and not self.stop_flag: retry_count += 1 print(f"[{self.name}] 强制重连中... (第{retry_count}/{max_retries}次尝试)") self.reconnect() if not self.get_live_state(): print(f"[{self.name}] 重连失败,5秒后重试...") time.sleep(5) if not self.get_live_state() and not self.stop_flag: print(f"[{self.name}] 达到最大重试次数,停止重连,监控线程将继续尝试") def _monitor_connection(self): """连接监控线程,定期检查连接状态并自动重连(新增退出机制)""" while not self.stop_flag: try: # 每5秒检查一次连接状态 time.sleep(5) # 检查是否需要重连 if not self.get_live_state() and not self.stop_flag: print(f"[{self.name}] 连接断开,开始重连...") self.force_reconnect() except Exception as e: print(f"[{self.name}] 监控线程错误: {e}") time.sleep(5) def _parse_state(self, data: bytes) -> Dict[str, Any]: """解析屏幕状态(需根据实际协议实现)""" return {} def _parse_params(self, data: bytes) -> Dict[str, Any]: """解析屏幕参数(需根据实际协议实现)""" return {} def display(self, color: int, run_mode: int, display_mode: int, is_program: int, speed: int, text: str): """核心显示方法""" if not self.get_live_state(): return self.send_ram(color, run_mode, display_mode, is_program, speed, text) def correct(self): """校正屏幕时间""" if not self.get_live_state(): return now = datetime.datetime.now() cmd = BxCmdFactory.NewBxCmdSystemClockCorrect(now) pack = NewBxDataPackCmd(cmd,0x0001) self.send(pack.Pack()) def format_speed(self, speed: int) -> str: """格式化速度参数""" if speed < 0 or speed > 99: return "00" return f"{speed:02d}" def send_ram(self, color: int, run_mode: int, display_mode: int, is_program: int, speed: int, text: str): """发送动态区节目""" # if is_program == 0: # ff = FlashFile() # ff.set_msg(self.format_speed(speed), color) # ff.set_mode(run_mode, display_mode) # ff.set_origin(0, True, 32) # ff.set_area(128, True, 96) # self.text_ram(ff, False) # elif is_program == 2: # ff = FlashFile() # ff.set_msg(text, color) # ff.set_mode(run_mode, display_mode) # ff.set_origin(0, True, 0) # ff.set_area(128, True, 32) # self.text_ram(ff, True) return def text_ram(self, ff: FlashFile, is_number: bool): """发送动态区节目实现(新增:车牌长度校验+绿牌FE编码区分)""" if not self.get_live_state(): return areas = [] # 编码转换 encoder = codecs.getencoder("gb2312") bytes_data = b"" id = 0 alignment = 0x00 disp_mode = ff.disp_mode dst_addr = 0x0003 if ff.color == Color.DEFAULT and not is_number: bytes_data, _ = encoder(ff.msg) elif is_number: # 完整车牌(如湘F65354) full_plate = ff.msg.strip() print(f"车牌号: {full_plate}") # ========== 核心新增:车牌合法性校验(按长度) ========== # 校验规则: # - 普通燃油车牌:7位(如湘F65354)→ FE000 # - 新能源绿色车牌:8位(如湘F653548)→ FE001 # - 非7/8位直接判定为非法,跳过处理 full_plate_len = len(full_plate) if full_plate_len not in [7, 8]: print(f"【校验失败】车牌{full_plate}长度{full_plate_len}位,非法(仅支持7/8位),跳过") return # 解析省份+号码(复用原有方法) province, number = self.parse_plate_utf8(full_plate) # 校验解析后的格式(省份非空 + 号码长度6/7位 + 号码首字符是字母) if not province or len(number) not in [6, 7] or not number[:1].isalpha() or not number[1:].isalnum(): print(f"【校验失败】车牌{full_plate}格式错误(解析:省份{province},号码{number}),跳过") return # ========== 原有逻辑:判断显示模式 ========== if len(number) == 6: disp_mode = DisplayMode.QUICK_PUNCH # ========== 核心修改:按长度判断绿牌,设置FE编码 ========== # 8位车牌(新能源绿牌)→ FE001,7位车牌(普通牌)→ FE000 fe_code = "001" if full_plate_len == 8 else "000" # 构建显示指令(替换FE编码,ff.color为屏幕颜色,保持不变) msg = f"\\FO000\\C{ff.color}{province}\\FE{fe_code}\\C{Color.GREEN}{number}" bytes_data, _ = encoder(msg) id = 1 dst_addr = 0x0001 elif not is_number: # 优化:使用与数字显示相同的格式,避免设备拒绝写入 msg = f"\\FO000\\C{ff.color}{ff.msg}" # print(ff.msg) bytes_data, _ = encoder(msg) alignment = 0x00 dst_addr = 0x0003 # 构建动态区 area = NewBxAreaProgram( id, ff.run_mode, disp_mode, alignment, ff.origin_x, ff.origin_y, ff.width, ff.height, bytes_data, False ) areas.append(area) # 打包指令并发送 cmd = BxCmdFactory.NewBxCmdSendDynamicArea(areas) pack = NewBxDataPackCmd(cmd, dst_addr) pack.SetDisplayType(1) data = pack.Pack() # print(f"打包后数据(十六进制): {data.hex()} {is_number}") self.send(data) if is_number: # 主屏幕:走队列 self.send(data) else: # 雷达屏幕:直接发送 send_success = self.send_direct(data) if not send_success: return # 读取响应 resp = self.read_resp() # resp_time_str = datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S.%f")[:-3] # 获取当前时间 毫秒级 # print(f"[响应时间 {resp_time_str}]") if not resp.IsAck(): remote_ip, remote_port = self.conn.getpeername() error_msg = resp.Error() or f"错误码: {resp.Err}" print( f"[{self.name}] 设备拒绝写文件: {resp.Error().ErrorCode} 错误:{remote_ip} {remote_port} {data.hex()}") return self.state_info["dyna_area_num"] = self.state_info.get("dyna_area_num", 0) + 1 def parse_plate_utf8(self, raw_plate: str) -> Tuple[str, str]: """解析车牌号(省份+号码)""" if not raw_plate: return "", "" province = raw_plate[0] number = raw_plate[1:] return province, number def text_flash(self, ft_list: List[FlashFile], name: str, is_number: bool, is_logo: bool): """发送静态文件节目(掉电保存)""" if not self.get_live_state(): return if is_logo: name = "LOGO" encoder = codecs.getencoder("gb18030") areas = [] for ff in ft_list: bytes_data = b"" if ff.color == Color.DEFAULT and not is_number: bytes_data, _ = encoder(ff.msg) elif is_number: province, number = self.parse_plate_utf8(ff.msg) print(f"车牌号: {ff.msg}, 省份: {province}, 号码: {number}") msg = f"\\FO000\\C{ff.color}{province}\\FE001\\C{Color.GREEN}{number}" bytes_data, _ = encoder(msg) elif not is_number: msg = f"\\FO000\\C{ff.color}{ff.msg}" bytes_data, _ = encoder(msg) area = NewBxAreaProgram( 0xff, ff.run_mode, ff.disp_mode, 0x00, ff.origin_x, ff.origin_y, ff.width, ff.height, bytes_data, False ) areas.append(area) # 构建文件并发送 file = NewBxFile(name, "", areas) cmd = file.NewCmdWriteFile() pack = NewBxDataPackCmd(cmd,0x0001) data = pack.Pack() self.send(data) resp = self.read_resp() if not resp.IsAck(): print(f"[{self.name}] 设备拒绝写文件: {resp.error}") return def read_resp(self) -> BxResp: """读取屏幕响应(优化:打印完整错误信息)""" resp = BxResp() if not self.get_live_state() or self.conn is None: return resp try: recv_data = self.conn.recv(1024) if not recv_data: print(f"[{self.name}] 未收到响应") return resp resp = resp.Parse(recv_data, len(recv_data)) return resp except socket.timeout: print(f"[{self.name}] 读取响应超时") return resp except Exception as e: print(f"[{self.name}] 读取响应失败: {e}") import traceback traceback.print_exc() return resp def state(self) -> Optional[BxResp]: """获取屏幕状态""" if not self.get_live_state(): return None cmd = BxCmdFactory.NewCmdState() pack = NewBxDataPackCmd(cmd,0x0001) self.send(pack.Pack()) return self.read_resp() def param(self) -> Optional[BxResp]: """获取屏幕参数""" if not self.get_live_state(): return None cmd = BxCmdFactory.NewCmdReadParams() pack = NewBxDataPackCmd(cmd,0x0001) self.send(pack.Pack()) return self.read_resp() def turn_on_off(self, on_off: bool): """开关屏""" if not self.get_live_state(): return cmd = BxCmdFactory.NewBxCmdTurnOnOff(on_off) pack = NewBxDataPackCmd(cmd,0x0001) self.send(pack.Pack()) def brightness(self, brightness_type: int, current_brightness: int, brightness_value: bytes): """设置亮度""" if not self.get_live_state(): return cmd = BxCmdFactory.NewCmdBrightness(brightness_type, current_brightness, brightness_value) pack = NewBxDataPackCmd(cmd,0x0001) print(f"亮度设置指令: {codecs.encode(pack.Pack(), 'hex')}") self.send(pack.Pack()) resp = self.read_resp() if not resp.IsAck(): print(f"[{self.name}] 设备拒绝设置亮度: {resp.error}") def timing_switch(self, on_off_set: List[Tuple[int, int]]): """设置定时开关屏""" if not self.get_live_state(): return cmd = BxCmdFactory.NewCmdTimingSwitch(on_off_set) pack = NewBxDataPackCmd(cmd,0x0001) self.send(pack.Pack()) def cancel_timing_switch(self): """取消定时开关屏""" if not self.get_live_state(): return cmd = BxCmdFactory.NewCmdCancelTimingSwitch() pack = NewBxDataPackCmd(cmd,0x0001) self.send(pack.Pack()) def del_file(self, del_files: List[str]): """删除静态文件""" if not self.get_live_state(): return cmd = BxCmdFactory.NewCmdDeleteFile(del_files) pack = NewBxDataPackCmd(cmd,0x0001) self.send(pack.Pack()) self.read_resp() if not del_files: self.state_info["program_num"] = 0 else: self.state_info["program_num"] = self.state_info.get("program_num", 0) - 1 def del_ram_text(self, *numbers: int): """ 对应 Go 的 DelRamText 方法:删除 RAM 文本动态区域 :param numbers: 可变参数,字节类型的动态区域编号 """ # 1. 检查屏幕是否处于活跃状态 if not self.get_live_state(): return # 2. 创建删除动态区域的指令(对应 Go 的 bx.NewCmdDelDynamicArea) cmd = BxCmdFactory.NewCmdDelDynamicArea(numbers) # 3. 打包指令(对应 Go 的 bx.NewBxDataPackCmd) pack = NewBxDataPackCmd(cmd,0x0003) # 4. 发送打包后的指令 data = pack.Pack() # print(f"打包后数据(十六进制): {data.hex()}") self.send(data) # 5. 读取响应 self.read_resp() # 6. 更新动态区域数量状态 if len(numbers) == 0: self.state_info["dyna_area_num"] = 0 else: self.state_info["dyna_area_num"] -= 1