295 lines
13 KiB
Python
295 lines
13 KiB
Python
|
|
# 注释掉/删除Linux的shebang行,替换为Windows兼容的写法(可选)
|
|||
|
|
# #!/usr/bin/env python3
|
|||
|
|
# -*- coding: utf-8 -*-
|
|||
|
|
"""
|
|||
|
|
故障注入平台UDP模拟客户端
|
|||
|
|
功能:向中间件多UDP端口发送16进制流数据,模拟故障注入平台的发包行为
|
|||
|
|
作者:工程测试专用
|
|||
|
|
日期:2026-02-26
|
|||
|
|
"""
|
|||
|
|
"""
|
|||
|
|
该程序专门模拟故障注入平台的 UDP 数据发送行为,向中间件的多 UDP 端口发送符合协议格式的 16 进制流数据,可验证中间件udp_handler.py的多端口非阻塞接收、数据解析、端口隔离等核心功能,支持单端口定向发送、多端口并发发送、自定义数据格式,全程适配工程的 16 进制流数据规范。
|
|||
|
|
一、测试程序核心功能
|
|||
|
|
模拟故障注入平台向中间件指定 UDP 端口发送 16 进制流数据;
|
|||
|
|
支持单端口循环发送、多端口并发发送(验证端口隔离性);
|
|||
|
|
可自定义测试数据(原始二进制 / 16 进制字符串),自动适配工程数据格式;
|
|||
|
|
实时打印发送日志,包含目标端口、数据长度、发送状态;
|
|||
|
|
支持按协议类型(UART/CAN/1553B/AD/OC)批量发送对应外设数据。
|
|||
|
|
"""
|
|||
|
|
"""
|
|||
|
|
使用步骤
|
|||
|
|
步骤 1:配置适配
|
|||
|
|
修改DEVICE_UDP_PORTS字典,确保与中间件device_config.csv中的外设名称 + UDP 端口完全一致;
|
|||
|
|
确认MIDDLEWARE_IP为中间件运行的 IP(本地测试填127.0.0.1,跨机器测试填中间件所在机器 IP);
|
|||
|
|
调整TEST_CONFIG中的参数:
|
|||
|
|
send_interval:发送间隔(秒),默认 1ms,模拟实时数据;
|
|||
|
|
concurrent_threads:并发线程数,默认 5,可根据机器性能调整;
|
|||
|
|
send_times_per_port:每个端口发送次数,默认 1000 次。
|
|||
|
|
步骤 2:运行准备
|
|||
|
|
先启动中间件主程序main.py,确保:
|
|||
|
|
配置初始化成功(日志显示 “加载 X 个外设”);
|
|||
|
|
UDP 服务初始化完成(日志显示 “监听端口:[8880,8881,...]”);
|
|||
|
|
安装依赖(仅需 loguru):
|
|||
|
|
bash
|
|||
|
|
运行
|
|||
|
|
pip install loguru
|
|||
|
|
步骤 3:执行测试
|
|||
|
|
将测试程序保存为fault_inject_simulator.py,放在工程根目录;
|
|||
|
|
运行测试程序:
|
|||
|
|
bash
|
|||
|
|
运行
|
|||
|
|
python fault_inject_simulator.py
|
|||
|
|
可选测试模式:
|
|||
|
|
单端口测试:取消send_single_port_cycle注释,测试指定外设端口的接收能力;
|
|||
|
|
多端口并发测试:取消send_multi_port_concurrent注释,测试多端口隔离性和并发接收能力。
|
|||
|
|
"""
|
|||
|
|
import sys
|
|||
|
|
import os
|
|||
|
|
import io
|
|||
|
|
|
|||
|
|
# 设置标准输出和标准错误的编码为UTF-8
|
|||
|
|
sys.stdout = io.TextIOWrapper(sys.stdout.buffer, encoding='utf-8')
|
|||
|
|
sys.stderr = io.TextIOWrapper(sys.stderr.buffer, encoding='utf-8')
|
|||
|
|
|
|||
|
|
import socket
|
|||
|
|
import binascii
|
|||
|
|
import time
|
|||
|
|
import threading
|
|||
|
|
import random
|
|||
|
|
from loguru import logger
|
|||
|
|
|
|||
|
|
from pathlib import Path
|
|||
|
|
# 获取当前文件的目录
|
|||
|
|
current_dir = os.path.dirname(os.path.abspath(__file__))
|
|||
|
|
# 获取上级目录(项目根目录)
|
|||
|
|
parent_dir = os.path.dirname(current_dir)
|
|||
|
|
# 将项目根目录添加到系统路径
|
|||
|
|
sys.path.append(parent_dir)
|
|||
|
|
#sys.path.append(str(Path(__file__).parent.parent.parent))
|
|||
|
|
from midware.config.base_config import UDP_CONFIG
|
|||
|
|
|
|||
|
|
# ==================== 配置项(与中间件base_config.py保持一致)====================
|
|||
|
|
# 中间件本地IP(故障注入平台访问的IP)
|
|||
|
|
MIDDLEWARE_IP = "127.0.0.1"
|
|||
|
|
# 故障注入平台接收端口(中间件回传数据的端口,可选)
|
|||
|
|
FAULT_RECV_PORT = 8889
|
|||
|
|
# 中间件监听的外设UDP端口列表(从device_config.csv复制)
|
|||
|
|
DEVICE_UDP_PORTS = {
|
|||
|
|
"Uart0": 4000,
|
|||
|
|
"CAN0": 4001,#": 8881,
|
|||
|
|
# "1553B0": 8882,
|
|||
|
|
# "AD0": 8883,
|
|||
|
|
# "OC0": 8884,
|
|||
|
|
# "Uart1": 8885,
|
|||
|
|
# "CAN2": 8886,
|
|||
|
|
# "1553B1": 8887,
|
|||
|
|
# "AD1": 8888,
|
|||
|
|
# "OC1": 8889
|
|||
|
|
}
|
|||
|
|
# 测试数据配置
|
|||
|
|
TEST_CONFIG = {
|
|||
|
|
"buffer_size": 32, # UDP缓冲区大小,4096
|
|||
|
|
"send_interval": 0.1, # 单端口发送间隔(秒,1ms,模拟实时数据)
|
|||
|
|
"concurrent_threads": 5, # 并发发送线程数
|
|||
|
|
"send_times_per_port": 10 # 每个端口循环发送次数
|
|||
|
|
}
|
|||
|
|
|
|||
|
|
# ==================== 日志配置 ====================
|
|||
|
|
logger.add("fault_inject_simulator.log", level="INFO",
|
|||
|
|
format="{time:YYYY-MM-DD HH:mm:ss} | {level} | {message}",
|
|||
|
|
rotation="10MB", retention="3 days")
|
|||
|
|
|
|||
|
|
class FaultInjectUDPSimulator:
|
|||
|
|
"""故障注入平台UDP模拟客户端"""
|
|||
|
|
def __init__(self):
|
|||
|
|
# 创建非阻塞UDP发送套接字
|
|||
|
|
self.send_sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
|
|||
|
|
self.send_sock.setblocking(False)
|
|||
|
|
# 开启地址复用,避免端口占用
|
|||
|
|
self.send_sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
|||
|
|
# 发送统计
|
|||
|
|
self.send_stats = {
|
|||
|
|
"total_send": 0,
|
|||
|
|
"success_send": 0,
|
|||
|
|
"failed_send": 0,
|
|||
|
|
"port_stats": {} # 按端口统计:{port: {"success":0, "failed":0}}
|
|||
|
|
}
|
|||
|
|
# 初始化端口统计
|
|||
|
|
for port in DEVICE_UDP_PORTS.values():
|
|||
|
|
self.send_stats["port_stats"][port] = {"success": 0, "failed": 0}
|
|||
|
|
logger.info("故障注入平台UDP模拟器初始化完成")
|
|||
|
|
|
|||
|
|
def _raw_to_hex_str(self, raw_data: bytes) -> str:
|
|||
|
|
"""
|
|||
|
|
原始二进制数据转16进制字符串(符合中间件数据格式要求)
|
|||
|
|
:param raw_data: 原始二进制数据
|
|||
|
|
:return: 16进制字符串(如b"test" → "74657374")
|
|||
|
|
"""
|
|||
|
|
return binascii.hexlify(raw_data).decode("utf-8")
|
|||
|
|
|
|||
|
|
def send_data_to_port(self, port: int, raw_data: bytes, dev_name: str = "unknown"):
|
|||
|
|
"""
|
|||
|
|
向指定UDP端口发送数据(核心方法)
|
|||
|
|
:param port: 中间件监听的UDP端口
|
|||
|
|
:param raw_data: 原始二进制数据(自动转为16进制流)
|
|||
|
|
:param dev_name: 外设名称(用于日志)
|
|||
|
|
"""
|
|||
|
|
try:
|
|||
|
|
# 转换为16进制字符串(中间件要求的格式)
|
|||
|
|
hex_data = self._raw_to_hex_str(raw_data)
|
|||
|
|
# 非阻塞发送数据
|
|||
|
|
self.send_sock.sendto(hex_data.encode("utf-8"), (MIDDLEWARE_IP, port))
|
|||
|
|
# 更新统计
|
|||
|
|
self.send_stats["total_send"] += 1
|
|||
|
|
self.send_stats["success_send"] += 1
|
|||
|
|
self.send_stats["port_stats"][port]["success"] += 1
|
|||
|
|
logger.debug(
|
|||
|
|
f"发送成功 | 外设:{dev_name} | 端口:{port} | "
|
|||
|
|
f"原始数据长度:{len(raw_data)}字节 | 16进制流:{hex_data[:32]}..."
|
|||
|
|
)
|
|||
|
|
return True
|
|||
|
|
except BlockingIOError:
|
|||
|
|
# 非阻塞发送缓冲区满
|
|||
|
|
self.send_stats["total_send"] += 1
|
|||
|
|
self.send_stats["failed_send"] += 1
|
|||
|
|
self.send_stats["port_stats"][port]["failed"] += 1
|
|||
|
|
logger.warning(f"发送失败 | 外设:{dev_name} | 端口:{port} | 原因:发送缓冲区满")
|
|||
|
|
return False
|
|||
|
|
except Exception as e:
|
|||
|
|
self.send_stats["total_send"] += 1
|
|||
|
|
self.send_stats["failed_send"] += 1
|
|||
|
|
self.send_stats["port_stats"][port]["failed"] += 1
|
|||
|
|
logger.error(
|
|||
|
|
f"发送异常 | 外设:{dev_name} | 端口:{port} | 错误:{str(e)}",
|
|||
|
|
exc_info=True
|
|||
|
|
)
|
|||
|
|
return False
|
|||
|
|
|
|||
|
|
def send_single_port_cycle(self, dev_name: str, send_times: int = None):
|
|||
|
|
"""
|
|||
|
|
单端口循环发送测试数据
|
|||
|
|
:param dev_name: 外设名称(如"Uart0",对应DEVICE_UDP_PORTS中的key)
|
|||
|
|
:param send_times: 发送次数,默认使用TEST_CONFIG中的配置
|
|||
|
|
"""
|
|||
|
|
if dev_name not in DEVICE_UDP_PORTS:
|
|||
|
|
logger.error(f"外设{dev_name}不存在,可选外设:{list(DEVICE_UDP_PORTS.keys())}")
|
|||
|
|
return
|
|||
|
|
port = DEVICE_UDP_PORTS[dev_name]
|
|||
|
|
send_times = send_times or TEST_CONFIG["send_times_per_port"]
|
|||
|
|
logger.info(f"开始单端口循环发送 | 外设:{dev_name} | 端口:{port} | 发送次数:{send_times}")
|
|||
|
|
|
|||
|
|
# 按协议类型生成测试数据(模拟真实外设数据)
|
|||
|
|
if "UART" in dev_name:
|
|||
|
|
raw_data = b"UART_DATA_" + str(random.randint(1000, 9999)).encode() + b"\x9c\x47" # 带CRC16
|
|||
|
|
elif "CAN" in dev_name:
|
|||
|
|
raw_data = b"CAN_DATA_" + str(random.randint(1000, 9999)).encode() + b"\x1a\x2b" # 带CRC16
|
|||
|
|
elif "1553B" in dev_name:
|
|||
|
|
raw_data = b"1553B_DATA_" + str(random.randint(1000, 9999)).encode() + b"\x0f\x3d" # 带CRC16
|
|||
|
|
elif "AD" in dev_name:
|
|||
|
|
raw_data = b"AD_DATA_" + str(random.randint(1000, 9999)).encode() + b"\x5f" # 带累加和
|
|||
|
|
elif "OC" in dev_name:
|
|||
|
|
raw_data = b"OC_DATA_" + str(random.randint(1000, 9999)).encode() + b"\x6a" # 带累加和
|
|||
|
|
else:
|
|||
|
|
raw_data = b"TEST_DATA_" + str(random.randint(1000, 9999)).encode()
|
|||
|
|
|
|||
|
|
# 循环发送
|
|||
|
|
for i in range(send_times):
|
|||
|
|
self.send_data_to_port(port, raw_data, dev_name)
|
|||
|
|
time.sleep(TEST_CONFIG["send_interval"])
|
|||
|
|
logger.info(f"单端口发送完成 | 外设:{dev_name} | 端口:{port} | 发送统计:{self.send_stats['port_stats'][port]}")
|
|||
|
|
|
|||
|
|
def send_multi_port_concurrent(self, dev_names: list = None, send_times: int = None):
|
|||
|
|
"""
|
|||
|
|
多端口并发发送测试(验证端口隔离性)
|
|||
|
|
:param dev_names: 外设名称列表,默认发送所有外设
|
|||
|
|
:param send_times: 每个端口发送次数,默认使用TEST_CONFIG中的配置
|
|||
|
|
"""
|
|||
|
|
dev_names = dev_names or list(DEVICE_UDP_PORTS.keys())
|
|||
|
|
send_times = send_times or TEST_CONFIG["send_times_per_port"]
|
|||
|
|
logger.info(f"开始多端口并发发送 | 外设列表:{dev_names} | 每端口发送次数:{send_times}")
|
|||
|
|
|
|||
|
|
# 定义单个端口的发送任务
|
|||
|
|
def send_task(dev_name):
|
|||
|
|
port = DEVICE_UDP_PORTS[dev_name]
|
|||
|
|
# 生成对应协议的测试数据
|
|||
|
|
if "UART" in dev_name:
|
|||
|
|
raw_data = b"UART_CONCURRENT_" + dev_name.encode() + b"_" + str(random.randint(1000, 9999)).encode()
|
|||
|
|
elif "CAN" in dev_name:
|
|||
|
|
raw_data = b"CAN_CONCURRENT_" + dev_name.encode() + b"_" + str(random.randint(1000, 9999)).encode()
|
|||
|
|
elif "1553B" in dev_name:
|
|||
|
|
raw_data = b"1553B_CONCURRENT_" + dev_name.encode() + b"_" + str(random.randint(1000, 9999)).encode()
|
|||
|
|
else:
|
|||
|
|
raw_data = b"CONCURRENT_" + dev_name.encode() + b"_" + str(random.randint(1000, 9999)).encode()
|
|||
|
|
# 循环发送
|
|||
|
|
for i in range(send_times):
|
|||
|
|
self.send_data_to_port(port, raw_data, dev_name)
|
|||
|
|
time.sleep(TEST_CONFIG["send_interval"])
|
|||
|
|
|
|||
|
|
# 创建并发线程
|
|||
|
|
threads = []
|
|||
|
|
max_threads = TEST_CONFIG["concurrent_threads"]
|
|||
|
|
active_threads = 0
|
|||
|
|
for dev_name in dev_names:
|
|||
|
|
if dev_name not in DEVICE_UDP_PORTS:
|
|||
|
|
logger.warning(f"外设{dev_name}不存在,跳过")
|
|||
|
|
continue
|
|||
|
|
# 控制并发线程数,避免资源耗尽
|
|||
|
|
while active_threads >= max_threads:
|
|||
|
|
time.sleep(0.001)
|
|||
|
|
active_threads = sum(1 for t in threads if t.is_alive())
|
|||
|
|
# 启动线程
|
|||
|
|
t = threading.Thread(target=send_task, args=(dev_name,), daemon=True)
|
|||
|
|
threads.append(t)
|
|||
|
|
t.start()
|
|||
|
|
active_threads += 1
|
|||
|
|
logger.debug(f"启动并发发送线程 | 外设:{dev_name} | 线程ID:{t.ident}")
|
|||
|
|
|
|||
|
|
# 等待所有线程完成
|
|||
|
|
for t in threads:
|
|||
|
|
t.join()
|
|||
|
|
logger.info("多端口并发发送完成 | 全局统计:"
|
|||
|
|
f"总发送{self.send_stats['total_send']} | "
|
|||
|
|
f"成功{self.send_stats['success_send']} | "
|
|||
|
|
f"失败{self.send_stats['failed_send']}")
|
|||
|
|
|
|||
|
|
def print_send_stats(self):
|
|||
|
|
"""打印最终发送统计"""
|
|||
|
|
logger.info("="*50 + " 发送统计汇总 " + "="*50)
|
|||
|
|
logger.info(f"全局统计 | 总发送:{self.send_stats['total_send']} | 成功:{self.send_stats['success_send']} | 失败:{self.send_stats['failed_send']}")
|
|||
|
|
logger.info("端口统计:")
|
|||
|
|
for port, stats in self.send_stats["port_stats"].items():
|
|||
|
|
dev_name = [k for k, v in DEVICE_UDP_PORTS.items() if v == port][0]
|
|||
|
|
total = stats["success"] + stats["failed"]
|
|||
|
|
success_rate = (stats["success"] / total * 100) if total > 0 else 0
|
|||
|
|
logger.info(f" 外设:{dev_name:6s} | 端口:{port:5d} | 发送:{total:5d} | 成功:{stats['success']:5d} | 失败:{stats['failed']:5d} | 成功率:{success_rate:.2f}%")
|
|||
|
|
logger.info("="*100)
|
|||
|
|
|
|||
|
|
def close(self):
|
|||
|
|
"""关闭套接字,释放资源"""
|
|||
|
|
self.send_sock.close()
|
|||
|
|
logger.info("故障注入平台UDP模拟器已关闭,资源释放完成")
|
|||
|
|
|
|||
|
|
# ==================== 测试执行入口 ====================
|
|||
|
|
def main():
|
|||
|
|
# 初始化模拟器
|
|||
|
|
simulator = FaultInjectUDPSimulator()
|
|||
|
|
print("故障注入平台UDP模拟器已启动,按Ctrl+C终止测试")
|
|||
|
|
try:
|
|||
|
|
# 可选测试模式:二选一或全选
|
|||
|
|
# 模式1:单端口循环发送(测试Uart0)
|
|||
|
|
# simulator.send_single_port_cycle("Uart0", send_times=50)
|
|||
|
|
|
|||
|
|
# 模式2:多端口并发发送(测试所有外设)
|
|||
|
|
simulator.send_multi_port_concurrent(dev_names=["Uart0", "CAN0", "1553B0", "AD0", "OC0"], send_times=20)
|
|||
|
|
|
|||
|
|
# 打印统计结果
|
|||
|
|
simulator.print_send_stats()
|
|||
|
|
except KeyboardInterrupt:
|
|||
|
|
logger.info("用户终止测试,打印最终统计...")
|
|||
|
|
simulator.print_send_stats()
|
|||
|
|
finally:
|
|||
|
|
# 关闭模拟器
|
|||
|
|
simulator.close()
|
|||
|
|
|
|||
|
|
if __name__ == "__main__":
|
|||
|
|
main()
|