Files
virtual_simulation_midware/tests/test_udp_connect_2.py

295 lines
13 KiB
Python
Raw Permalink Normal View History

2026-06-16 15:40:19 +08:00
# 注释掉/删除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()