 的PLC-HMI工程项目(十四)后续完善和改进:上位机数据解析的防粘包、残包、废包、拆包)
基于基于socket的 TCP上位机通信内容解析程序-CSDN博客改进后的PLC.py:import socket import struct from typing import NamedTuple import numpy as np from PySide6.QtCore import QObject, Signal, QTimer, Slot, QMutex, QThread, QMutexLocker from .settings import * # 字节长度的字典 type_bytes {Byte: 1,Bool: 1, Int: 2, Real: 4, Word:2, DInt:4, DWord:4, LInt:4, LWord:4, LReal:8} def set_bit(b: bytes, bit_idx: int, value: int) - bytes: 修改 bytes 变量的某一位 :param b: 原始字节 :param bit_idx: 要修改的位下标0~70 是最低位 :param value: 目标值 0 或 1 :return: 修改后的新 bytes # 转成可变的 bytearray arr bytearray(b) # 取出目标字节 byte_val arr[0] if value 1: # 置 1按位或 | # b 1 bit_idx # byte_val byte_val | b byte_val | 1 bit_idx else: # 置 0按位与 取反 ~ byte_val ~(1 bit_idx) # 写回 arr[0] byte_val # 转回 bytes return bytes(arr) def get_bit(byte_val: bytes, bit_idx: int) - int: 判断一个字节的某一位是 0 还是 1 :param byte_val: 输入字节只判断第一个字节 :param bit_idx: 位索引0~70最低位 :return: 0 或 1 # return 1 if (byte_val[0] (1 bit_idx)) ! 0 else 0 return True if (byte_val[0] (1 bit_idx)) ! 0 else False def crc16_modbus(data: bytes) - int: crc 0xFFFF for byte in data: crc ^ byte for _ in range(8): if crc 1: crc (crc 1) ^ 0xA001 else: crc 1 return crc def hex_str(data: bytes) - str: 字节转空格分隔的hex字符串 return .join(f{b:02X} for b in data) class VariableFrameParser: HEADER b\x41\x41 TAIL b\x42\x42 CRC_LEN 2 MAX_FRAME 1024 # 单帧最大字节数保护根据实际调整 def __init__(self, debug: bool True): self.buffer bytearray() self.debug debug self.frame_count 0 def log(self, msg: str): if self.debug: print(f [解析] {msg}) def feed(self, data: bytes): 喂入任意长度TCP数据返回解析出的完整帧列表 if not data: return [] self.log(f喂入 {len(data)} 字节: {hex_str(data)}) self.buffer.extend(data) frames [] while True: frame self._try_parse_one() if frame is None: break frames.append(frame) if not frames and self.debug: self.log(f缓冲区现有 {len(self.buffer)} 字节暂未形成完整帧) return frames def _try_parse_one(self): buf self.buffer # 外层循环: 支持跳帧后重新同步从新帧头开始 while True: # ── 第1步: 找帧头 ── header_pos buf.find(self.HEADER) print(header_pos:, header_pos) if header_pos -1: # 没找到帧头保留最后1字节可能是0x58下一包凑第二个0x58 if len(buf) 1: discarded bytes(buf[:-1]) del buf[:-1] return None if header_pos 0: discarded bytes(buf[:header_pos]) self.log(f帧头前有 {header_pos} 字节垃圾: {hex_str(discarded)}已丢弃) del buf[:header_pos] # 丢弃帧头前的垃圾数据 # ── 第2步: 逐个尝试帧尾用CRC验证真假 ── search_pos 2 attempt 0 while True: tail_pos buf.find(self.TAIL, search_pos) print(tail_pos:, tail_pos) if tail_pos -1: # 缓冲区中所有帧尾 59 59 都已尝试完毕全部 CRC 失败 # 检查是否存在下一个帧头且其后有完整帧结构 next_header buf.find(self.HEADER, 2) if next_header ! -1: next_tail buf.find(self.TAIL, next_header 2) if next_tail ! -1 and len(buf) next_tail 2 self.CRC_LEN: # 判定当前帧头为错误帧跳过至新帧头 del buf[:next_header] break # 跳出内层回到外层重新找帧头 # 没有下一个有效帧头可能是正常拆包等更多数据 if len(buf) self.MAX_FRAME: # 缓冲区超过{self.MAX_FRAME}字节仍无有效帧判定当前帧头为假跳过1字节重新同步 del buf[:1] break # 跳出内层回到外层重新找帧头 # 未找到有效帧尾等待更多数据 return None # 单帧长度保护 if tail_pos 2 self.CRC_LEN self.MAX_FRAME: # 帧尾候选{tail_pos}超出最大帧长当前帧头无效跳过1字节 del buf[:1] break # 跳出内层回到外层重新找帧头 attempt 1 frame_end tail_pos 2 self.CRC_LEN if len(buf) frame_end: # 帧尾候选{tail_pos}但CRC未到齐需等待 return None # ── 第3步: CRC校验 ── crc_range bytes(buf[:tail_pos 2]) # 帧头数据帧尾 crc_received struct.unpack(H, buf[tail_pos 2:frame_end])[0] crc_calc crc16_modbus(crc_range) if crc_calc crc_received: payload bytes(buf[2:tail_pos]) # 取出数据 raw bytes(buf[:frame_end]) # 原始字节 del buf[:frame_end] # 删除本帧数据 self.frame_count 1 print(f★ 解析成功! 有效数据({len(payload)}B): {hex_str(payload)}) # return {data: payload, raw: raw, seq: self.frame_count} return payload # CRC不通过 → 这个0x59 0x59是数据中的内容继续找下一个 search_pos tail_pos 1 # 区域变量 class AreaVar(NamedTuple): area: str # I、Q、M、DB DBnum: int # DB号I/Q/M0 offset: int # 字节偏移 byte_count: int # 字节数 value_bytes: bytes b # 数据的字节 def __eq__(self, other): return self.value_bytes other.value_bytes # UI变量 class UiVari(QObject): valueChanged Signal(object) def __init__(self, typ:str, byte_count:int, area:str, offset:int, DBnum0, bit0, valueNone ): UI变量用来与UI部件交互和触发通信操作 :param typ: :param typ: 类型 :param byte_count: 字节数 :param area: 区域 :param offset: 地址偏移 :param DBnum: DB编号 :param bit: 位 :param value: 值 super().__init__() self.typ typ self.byte_count byte_count self.area area self.DBnum DBnum self.offset offset self.bit bit self._value value def __eq__(self, other): return self.value other.value property def value(self): 获取变量值 return self._value value.setter def value(self, new_val): 设置值变化时自动发射信号 if self._value ! new_val: self._value new_val self.valueChanged.emit(new_val) def set_value(self, val): 外部设置值接口语义更清晰 self.value val # 上行报文变量 class UplinkMsg(QObject): valueChanged Signal(object) def __init__(self): super().__init__() self._value b property def value(self): 获取变量值 return self._value value.setter def value(self, new_val): 设置值变化时自动发射信号 if self._value ! new_val: self._value new_val self.valueChanged.emit(new_val) def set_value(self, val): 外部设置值接口语义更清晰 self.value val # Numpy格式数据区存储支持 I/Q/M/DB class DataStore: def __init__(self, max_db_number: int 2048, area_size: int 2048): 统一数据存储I、Q、M、DB :param max_db_number: 最大 DB 块号 :param area_size: 每个区域默认大小字节 self.area_size area_size # I / Q / M 区一维数组DBnum0 self.I_area np.zeros(area_size, dtypenp.uint8) # I区 self.Q_area np.zeros(area_size, dtypenp.uint8) # Q区 self.M_area np.zeros(area_size, dtypenp.uint8) # M区 # DB区二维数组 [DB号][偏移] self.DB_area np.zeros((max_db_number 1, area_size), dtypenp.uint8) # 变量表管理器按变量名直接读写 class VariManager: def __init__(self, store: DataStore): 变量表管理器按变量名直接读写 :param store: 变量存储区 self.store store self.vari_lock QMutex() self.store_lock QMutex() # UiVari变量转为AreaVar变量 def UiVari2AreaVar(self, ui_vari:UiVari): # def pack_value(self, ui_vari:UiVari, value): typ ui_vari.typ # 变量格式 value_bytes b # 先构造一个AreaVar变量 area_var AreaVar( areaui_vari.area, DBnumui_vari.DBnum, offsetui_vari.offset, byte_countui_vari.byte_count ) if ui_vari.value is not None: value ui_vari.value # 打包成 bytes if typ Bool: value_bytes self.read_block(area_var) bit ui_vari.bit value_bytes set_bit(value_bytes, bit, value) elif typ Byte: value_bytes value elif typ Int: value_bytes struct.pack(h, value) # 大端 Int16 elif typ Real: value_bytes struct.pack(f, value) # 大端 Float32 elif typ Word or typ DWord or typ LWord: value_bytes value elif typ LInt: value_bytes struct.pack(i, value) elif typ LReal: value_bytes struct.pack(d, value) else: raise TypeError(f不支持类型: {typ}) area_var AreaVar( areaui_vari.area, DBnumui_vari.DBnum, offsetui_vari.offset, byte_countui_vari.byte_count, value_bytesvalue_bytes ) return area_var # ------------------- 按变量 写值 ------------------- def write_vari(self, ui_vari:UiVari, value): ui_vari.set_value(value) with QMutexLocker(self.vari_lock): area_var self.UiVari2AreaVar(ui_vari) self.write_block(area_var) # ------------------- 按变量 读值 ------------------- def read_vari(self, ui_vari:UiVari): with QMutexLocker(self.vari_lock): typ ui_vari.typ # 变量格式 # 先构造一个AreaVar变量 area_var self.UiVari2AreaVar(ui_vari) # 读取变量字节串 data self.read_block(area_var) # 解析值 if typ Bool: # value 1 if get_bit(data, ui_vari.bit) True else 0 value get_bit(data, ui_vari.bit) elif typ Byte: value data elif typ Int: value struct.unpack(h, data)[0] elif typ Real: value struct.unpack(f, data)[0] elif typ Word or typ DWord or typ LWord: value data elif typ LInt: value struct.unpack(i, data)[0] elif typ LReal: value struct.unpack(d, data)[0] else: raise TypeError(f不支持类型: {typ}) # 将变量本身的value也修改 ui_vari.set_value(value) # print(ui_vari) return value # ------------------- 按块 写值 ------------------- def write_block(self, area_var: AreaVar): 统一写入接口支持 I/Q/M/DB with QMutexLocker(self.store_lock): data np.frombuffer(area_var.value_bytes, dtypenp.uint8) start area_var.offset end start area_var.byte_count if area_var.area I: self.store.I_area[start:end] data elif area_var.area Q: self.store.Q_area[start:end] data elif area_var.area M: self.store.M_area[start:end] data elif area_var.area DB: self.store.DB_area[area_var.DBnum, start:end] data else: raise ValueError(f不支持的区域: {area_var.area}) # ------------------- 按块 读值 ------------------- def read_block(self, area_var: AreaVar) - bytes: 统一读取接口支持 I/Q/M/DB with QMutexLocker(self.store_lock): start area_var.offset end start area_var.byte_count if area_var.area I: arr self.store.I_area[start:end] elif area_var.area Q: arr self.store.Q_area[start:end] elif area_var.area M: arr self.store.M_area[start:end] elif area_var.area DB: arr self.store.DB_area[area_var.DBnum, start:end] else: raise ValueError(f不支持的区域: {area_var.area}) return arr.tobytes() # 用于TCP通信的周期运行的子线程工作者 class ThreadWorker(QObject): _signal_start Signal() # 开始运行的信号 _signal_stop Signal() # 停止信号 _signal_run_now Signal() # 打破定时器周期立即运行一次 def __init__(self, target, looping, loop_interval0, *args, **kwargs): :param target: 需要在子线程中运行的目标函数 :param looping: 是周期运行还是一次性任务 :param loop_interval: 周期运行的间隙时间 :param args: 附带参数 :param kwargs: super().__init__() self._loop_timer None # 循环运行定时器 self.looping looping # 是否循环运行 self.loop_interval loop_interval # 循环运行间隔 self.target target # 目标函数 self.args args self.kwargs kwargs self.started False # 已经开始 self._signal_start.connect(self._start) self._signal_stop.connect(self._stop) self._signal_run_now.connect(self._run_target) # 内部接口在子线程内用槽函数操作定时器 Slot() def _start(self): self._run_target() if self.looping: if self._loop_timer is None: self._loop_timer QTimer() self._loop_timer.setInterval(self.loop_interval) self._loop_timer.timeout.connect(self._run_target) self._loop_timer.start() # 内部接口在子线程内用槽函数操作定时器 Slot() def _stop(self): if self._loop_timer is not None: self._loop_timer.stop() # 外部接口在调用线程内发射信号 def start(self): self.started True self._signal_start.emit() # 外部接口在调用线程内发射信号 def stop(self): self.started False self._signal_stop.emit() # 外部接口在调用线程内发射信号 def run_now(self): self._signal_run_now.emit() # 槽函数运行目标函数 Slot() def _run_target(self): if self._loop_timer is not None: self._loop_timer.stop() # 停止定时器防止信号堆积 if self.started: self.target(*self.args, **self.kwargs) if self._loop_timer is not None: self._loop_timer.start() # 重启定时器 # TCP客户端 class TcpClient(QObject): def __init__(self, ip, port): super().__init__() self.ip ip self.port port self.socket None self.is_connected False self.send_bytes b # 发送的字节 self.rsv_msg UplinkMsg() # 接收到的字节 self.rsv_lock QMutex() self.send_lock QMutex() ############## 创建线程和worker ################ # 连接服务器 self.connect_thread QThread() # 线程 self.connect_worker ThreadWorker(targetself.connect_server, loopingTrue, loop_intervalRECONNECT_INTERVAL) # worker self.connect_worker.moveToThread(self.connect_thread) # 把worker移动到线程中 self.connect_thread.start() # 启动线程 # 在子线程中周期发送数据 self.send_thread QThread() # 线程 self.send_worker ThreadWorker(targetself.send_data, loopingTrue, loop_intervalHEARTBEAT_INTERVAL) # worker self.send_worker.moveToThread(self.send_thread) # 把worker移动到线程中 self.send_thread.start() # 启动线程 # 在子线程中持续接收数据 self.recv_thread QThread() # 线程 self.recv_worker ThreadWorker(targetself.recv_data, loopingTrue, loop_interval0) # worker self.recv_worker.moveToThread(self.recv_thread) # 把worker移动到线程中 self.recv_thread.start() # 启动线程 self.uplink_parser VariableFrameParser() # 上行数据的解析程序 def connect_server(self): 建立 TCP 连接 if self.is_connected: return try: # 停止发送工作者 self.send_worker.stop() # 停止接收工作者 self.recv_worker.stop() # 创建 TCP socket self.socket socket.socket(socket.AF_INET, socket.SOCK_STREAM) self.socket.settimeout(CONNECT_TIMEOUT) self.socket.connect((self.ip, self.port)) # 连接服务器 self.is_connected True self.connect_worker.stop() # 连接完成停止连接worker print(f✅ 成功连接到服务器 {self.ip}:{self.port}) # 启动发送工作者 self.send_worker.start() # 启动接收工作者 self.recv_worker.start() except Exception as e: print(f❌ 连接失败{RECONNECT_INTERVAL}秒后重试... 错误{e}) self.is_connected False # self.connect_worker.stop() # self.send_worker.stop() # self.recv_worker.stop() def send_data(self): 发送数据 if not self.is_connected: print(⚠️ 未连接无法发送) return False try: with QMutexLocker(self.send_lock): # self.send_bytes # 发送数据 self.socket.send(self.send_bytes) # print(f 已发送{self.send_bytes}) self.send_bytes DOWNLINK_FRAME_HEAD DOWNLINK_FRAME_END # 发送完成后清空发送区 return True except Exception as e: print(f❌ 发送失败{e}) self.is_connected False self.send_worker.stop() self.recv_worker.stop() self.connect_worker.start() return False # 打破周期立即发送一个数据 def send_now(self): self.send_worker.run_now() def recv_data(self): 接收数据处理半包/粘包 if not self.is_connected: print(⚠️ 未连接无法接收) return False # buffer b # 清空接收缓存区 try: data self.socket.recv(1024) # print(data: , type(data)) if not data: print( 服务器断开连接) self.is_connected False self.send_worker.stop() self.recv_worker.stop() self.connect_worker.start() return False frames self.uplink_parser.feed(data) # 解析到的帧数据 # print(frames: , frames) with QMutexLocker(self.rsv_lock): for frame in frames: self.rsv_msg.set_value(frame) print(f 收到数据{frame}) except socket.timeout: return True except Exception as e: print(f❌ 接收异常{e}) self.is_connected False self.send_worker.stop() self.recv_worker.stop() self.connect_worker.start() def start(self): 启动客户端 print(f 启动 TCP 客户端目标{self.ip}:{self.port}) self.connect_worker.start() def quit(self): self.is_connected False self.connect_worker.stop() self.send_worker.stop() self.recv_worker.stop() self.connect_thread.quit() self.recv_thread.quit() self.send_thread.quit() # 把TcpClient和PlcPartner整合成一个然后定义一个处理上行帧的函数下行数据加进去CRC # PLC通信伙伴 class PlcPartner(QObject): def __init__(self, plc_ip, plc_port, vari_table): super().__init__() self.all_vari_list None self.vs DataStore() # 创建变量存储区 self.vm VariManager(self.vs) # 创建变量管理器 self.client TcpClient(plc_ip, plc_port) # 客户端 self.vari_table vari_table # 变量表 self.update_workers [] # 周期更新变量的工作者们 self.parse_uplink_worker None # 解析上行报文的工作者 # 写PLC变量 def write_vari(self, ui_vari: UiVari, value): self.vm.write_vari(ui_vari, value) # 先写入本地 # 如果是写一个单bit把位Int格式打包成字节 if ui_vari.typ Bool: bit ui_vari.bit bit_bytes struct.pack(i, bit) else: bit_bytes b\xff\xff\xff\xff # 如果不是写单个的bit位字节为 FFFFFFFF # 将Ui变量转为区域变量 area_var self.vm.UiVari2AreaVar(ui_vari) # 拼接下行报文 area_var_bytes self.AreaVar2msg(area_var) # 把区域变量的信息转成字节携带在下行报文中 send_bytes b.join([DOWNLINK_FRAME_HEAD, area_var_bytes, bit_bytes, area_var.value_bytes, DOWNLINK_FRAME_END]) # print(send_bytes) #生成CRC crc_bytes CRC16.new(send_bytes).crcValue.to_bytes(2, byteorderbig) # crc_bytes struct.pack(i, crc_result) # print(crc_bytes) with QMutexLocker(self.client.send_lock): #DOWNLINK_FRAME_HEAD DOWNLINK_FRAME_END self.client.send_bytes send_bytes crc_bytes # 发送报文到PLC self.send_now() def read_vari(self, ui_vari): # print(2) # with QMutexLocker(self.lock): v self.vm.read_vari(ui_vari) return v def write_block(self, area_var: AreaVar): # with QMutexLocker(self.lock): self.vm.write_block(area_var) def read_block(self, area_var: AreaVar) - bytes: # with QMutexLocker(self.lock): v self.vm.read_block(area_var) return v def start(self): self.client.start() # 启动客户端 self.start_updater() # 启动变量自动更新 def quit(self): self.client.quit() self.quit_updater() # 打破周期立即发送 def send_now(self): self.client.send_now() def AreaVar2msg(self, area_vari: AreaVar): # 把area变量的信息转换成字节 if area_vari.area I: area_bytes b\x81\x81 elif area_vari.area Q: area_bytes b\x82\x82 elif area_vari.area M: area_bytes b\x83\x83 elif area_vari.area DB: area_bytes b\x84\x84 else: raise ValueError(下行报文组态错误未识别的分区) # 把DB号转换成字节 db_num_bytes struct.pack(i, area_vari.DBnum) # 把offset转换成字节 offset_bytes struct.pack(i, area_vari.offset) # 把byte_count转换成字节 byte_count_bytes struct.pack(i, area_vari.byte_count) out_bytes b.join([area_bytes, db_num_bytes, offset_bytes, byte_count_bytes]) # print(out_bytes) return out_bytes def get_vari_update_table(self,varies: list): 从变量表的所有变量列表中获取到每个周期需要更新的变量内容 :param varies: :return: periods set() for v in varies: periods.add(v[1]) periods sorted(list(periods)) varis [] # 最终的输出结果 for p in periods: # 所有周期 # print(p) o [] for v in varies: if v[1] p: o.append(v[0]) varis.append(o) # print(periods, varis) return periods, varis # 开始更新变量 def start_updater(self): variables self.vari_table.variables all_period, self.all_vari_list self.get_vari_update_table(variables) for i, p in enumerate(all_period): worker ThreadWorker(targetself.update_vari, loopingTrue, loop_intervalp ) th QThread() worker.moveToThread(th) worker.start() th.start() self.update_workers.append([worker,th]) # 新建的更新worker和thread # 解析上行报文的工作者 self.parse_uplink_worker ThreadWorker(targetself.parse_uplink_msg, loopingFalse # 非定时器周期运行每次run_now()运行一次 ) th QThread() self.parse_uplink_worker.moveToThread(th) self.parse_uplink_worker.start() th.start() self.update_workers.append([self.parse_uplink_worker,th]) # 当上行报文的内容发生变化就解析一次 self.client.rsv_msg.valueChanged.connect(self.parse_uplink_worker.run_now) def quit_updater(self): for worker in self.update_workers: worker[0].stop() worker[1].quit() # 更新变量通过读取变量表中的变量和周期定义 def update_vari(self): for vari_list in self.all_vari_list: for vari in vari_list: self.read_vari(vari) # if value ! vari.value: # self.vm.write_vari(vari, value) Slot() def parse_uplink_msg(self): msg_bytes self.client.rsv_msg.value # print(msg_bytes, msg_bytes) # print(msg_bytes) if msg_bytes b: return if len(msg_bytes) 19: raise BufferError(msg_bytes, ❌数据长度不足) # head msg_bytes[0:2] # 帧头 # tail msg_bytes[-4:-2] # 帧尾 # if head ! UPLINK_FRAME_HEAD or tail ! UPLINK_FRAME_END: # raise BufferError(f❌数据校验失败帧头尾错误) # CRC # crc_int int.from_bytes(msg_bytes[-2:], byteorderbig) # crc_result CRC16.new(msg_bytes[:-2]).crcValue # if crc_result ! crc_int: # raise BufferError(f❌数据校验失败CRC失败) # # 截取上行数据的有效部分 up_data msg_bytes # lenxxx len(up_data) # print(fup_data: {lenxxx},up_data) pos 0 # 定义一个指针用来遍历字节 AreaVars [] # 输出的变量列表 # 解析数据 while pos len(up_data): try: area AREA_DICT[up_data[pos]] except KeyError: raise BufferError(f❌数据校验失败数据中区域代码错误) pos 2 DBnum int.from_bytes(up_data[pos:pos 4], byteorderbig) print(DBnum: , DBnum) pos 4 offset int.from_bytes(up_data[pos:pos 4], byteorderbig) print(offset: , offset) pos 4 byte_count int.from_bytes(up_data[pos:pos 4], byteorderbig) print(byte_count: , byte_count) pos 4 bytes_get up_data[pos:pos byte_count] print(bytes_get: , bytes_get) pos byte_count AreaVars.append(AreaVar(area, DBnum, offset, byte_count, bytes_get)) print(pos: ,pos) print(AreaVars) if pos ! len(up_data): raise BufferError(f❌数据校验失败数据长度错误解析失败) # with QMutexLocker(self.lock): for var in AreaVars: self.write_block(var)测试下位将DB1.DBB0--DBB4的字节内容“11 22 33 44 55”传输到上位并映射到上位数据池更新显示下位将油泵1的压力设定内容“10.0”传输到上位并映射到上位数据池更新显示