| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597 |
- 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
|