package cn.iocoder.yudao.module.iot.gateway.protocol.websocket;
|
|
import cn.hutool.core.map.MapUtil;
|
import cn.hutool.core.util.IdUtil;
|
import cn.hutool.core.util.StrUtil;
|
import cn.iocoder.yudao.module.iot.core.biz.dto.IotDeviceAuthReqDTO;
|
import cn.iocoder.yudao.module.iot.core.enums.IotDeviceMessageMethodEnum;
|
import cn.iocoder.yudao.module.iot.core.mq.message.IotDeviceMessage;
|
import cn.iocoder.yudao.module.iot.core.topic.event.IotDeviceEventPostReqDTO;
|
import cn.iocoder.yudao.module.iot.core.topic.property.IotDevicePropertyPostReqDTO;
|
import cn.iocoder.yudao.module.iot.core.util.IotDeviceAuthUtils;
|
import cn.iocoder.yudao.module.iot.gateway.serialize.IotMessageSerializer;
|
import cn.iocoder.yudao.module.iot.gateway.serialize.json.IotJsonSerializer;
|
import io.vertx.core.Vertx;
|
import io.vertx.core.http.WebSocket;
|
import io.vertx.core.http.WebSocketClient;
|
import io.vertx.core.http.WebSocketConnectOptions;
|
import lombok.extern.slf4j.Slf4j;
|
import org.junit.jupiter.api.AfterAll;
|
import org.junit.jupiter.api.BeforeAll;
|
import org.junit.jupiter.api.Disabled;
|
import org.junit.jupiter.api.Test;
|
|
import java.util.concurrent.CountDownLatch;
|
import java.util.concurrent.TimeUnit;
|
import java.util.concurrent.atomic.AtomicReference;
|
|
/**
|
* IoT 网关子设备 WebSocket 协议集成测试(手动测试)
|
*
|
* <p>测试场景:子设备(IotProductDeviceTypeEnum 的 SUB 类型)通过网关设备代理上报数据
|
*
|
* <p><b>重要说明:子设备无法直接连接平台,所有请求均由网关设备(Gateway)代为转发。</b>
|
*
|
* <p>使用步骤:
|
* <ol>
|
* <li>启动 yudao-module-iot-gateway 服务(WebSocket 端口 8094)</li>
|
* <li>确保子设备已通过 {@link IotGatewayDeviceWebSocketProtocolIntegrationTest#testTopoAdd()} 绑定到网关</li>
|
* <li>运行以下测试方法:
|
* <ul>
|
* <li>{@link #testAuth()} - 子设备认证</li>
|
* <li>{@link #testPropertyPost()} - 子设备属性上报(由网关代理转发)</li>
|
* <li>{@link #testEventPost()} - 子设备事件上报(由网关代理转发)</li>
|
* </ul>
|
* </li>
|
* </ol>
|
*
|
* <p>注意:WebSocket 协议是有状态的长连接,认证成功后同一连接上的后续请求无需再携带认证信息
|
*
|
* @author 芋道源码
|
*/
|
@Slf4j
|
@Disabled
|
public class IotGatewaySubDeviceWebSocketProtocolIntegrationTest {
|
|
private static final String SERVER_HOST = "127.0.0.1";
|
private static final int SERVER_PORT = 8094;
|
private static final String WS_PATH = "/ws";
|
private static final int TIMEOUT_SECONDS = 5;
|
|
private static Vertx vertx;
|
|
// ===================== 序列化器选择 =====================
|
|
private static final IotMessageSerializer SERIALIZER = new IotJsonSerializer();
|
|
// ===================== 网关子设备信息(根据实际情况修改,从 iot_device 表查询子设备) =====================
|
|
private static final String PRODUCT_KEY = "jAufEMTF1W6wnPhn";
|
private static final String DEVICE_NAME = "chazuo-it";
|
private static final String DEVICE_SECRET = "d46ef9b28ab14238b9c00a3a668032af";
|
|
@BeforeAll
|
public static void setUp() {
|
vertx = Vertx.vertx();
|
}
|
|
@AfterAll
|
public static void tearDown() {
|
if (vertx != null) {
|
vertx.close();
|
}
|
}
|
|
// ===================== 认证测试 =====================
|
|
/**
|
* 子设备认证测试
|
*/
|
@Test
|
public void testAuth() throws Exception {
|
// 1.1 构建认证消息
|
IotDeviceAuthReqDTO authInfo = IotDeviceAuthUtils.getAuthInfo(PRODUCT_KEY, DEVICE_NAME, DEVICE_SECRET);
|
IotDeviceAuthReqDTO authReqDTO = new IotDeviceAuthReqDTO()
|
.setClientId(authInfo.getClientId())
|
.setUsername(authInfo.getUsername())
|
.setPassword(authInfo.getPassword());
|
IotDeviceMessage request = IotDeviceMessage.of(IdUtil.fastSimpleUUID(), "auth", authReqDTO, null, null, null);
|
// 1.2 序列化
|
byte[] payload = SERIALIZER.serialize(request);
|
String jsonMessage = StrUtil.utf8Str(payload);
|
log.info("[testAuth][Serialize: {}, 请求消息: {}]", SERIALIZER.getType(), request);
|
|
// 2.1 创建 WebSocket 连接(同步)
|
WebSocket ws = createWebSocketConnection();
|
log.info("[testAuth][WebSocket 连接成功]");
|
|
// 2.2 发送并等待响应
|
String response = sendAndReceive(ws, jsonMessage);
|
|
// 3. 解码响应
|
if (response != null) {
|
IotDeviceMessage responseMessage = SERIALIZER.deserialize(StrUtil.utf8Bytes(response));
|
log.info("[testAuth][响应消息: {}]", responseMessage);
|
} else {
|
log.warn("[testAuth][未收到响应]");
|
}
|
|
// 4. 关闭连接
|
ws.close();
|
}
|
|
// ===================== 子设备属性上报测试 =====================
|
|
/**
|
* 子设备属性上报测试
|
*/
|
@Test
|
public void testPropertyPost() throws Exception {
|
// 1.1 创建 WebSocket 连接(同步)
|
WebSocket ws = createWebSocketConnection();
|
log.info("[testPropertyPost][WebSocket 连接成功]");
|
|
// 1.2 先进行认证
|
IotDeviceMessage authResponse = authenticate(ws);
|
log.info("[testPropertyPost][认证响应: {}]", authResponse);
|
log.info("[testPropertyPost][子设备属性上报 - 请求实际由 Gateway 代为转发]");
|
|
// 2.1 构建属性上报消息
|
IotDeviceMessage request = IotDeviceMessage.of(
|
IdUtil.fastSimpleUUID(),
|
IotDeviceMessageMethodEnum.PROPERTY_POST.getMethod(),
|
IotDevicePropertyPostReqDTO.of(MapUtil.<String, Object>builder()
|
.put("power", 100)
|
.put("status", "online")
|
.put("temperature", 36.5)
|
.build()),
|
null, null, null);
|
// 2.2 序列化
|
byte[] payload = SERIALIZER.serialize(request);
|
String jsonMessage = StrUtil.utf8Str(payload);
|
log.info("[testPropertyPost][Serialize: {}, 请求消息: {}]", SERIALIZER.getType(), request);
|
|
// 3.1 发送并等待响应
|
String response = sendAndReceive(ws, jsonMessage);
|
// 3.2 解码响应
|
if (response != null) {
|
IotDeviceMessage responseMessage = SERIALIZER.deserialize(StrUtil.utf8Bytes(response));
|
log.info("[testPropertyPost][响应消息: {}]", responseMessage);
|
} else {
|
log.warn("[testPropertyPost][未收到响应]");
|
}
|
|
// 4. 关闭连接
|
ws.close();
|
}
|
|
// ===================== 子设备事件上报测试 =====================
|
|
/**
|
* 子设备事件上报测试
|
*/
|
@Test
|
public void testEventPost() throws Exception {
|
// 1.1 创建 WebSocket 连接(同步)
|
WebSocket ws = createWebSocketConnection();
|
log.info("[testEventPost][WebSocket 连接成功]");
|
|
// 1.2 先进行认证
|
IotDeviceMessage authResponse = authenticate(ws);
|
log.info("[testEventPost][认证响应: {}]", authResponse);
|
log.info("[testEventPost][子设备事件上报 - 请求实际由 Gateway 代为转发]");
|
|
// 2.1 构建事件上报消息
|
IotDeviceMessage request = IotDeviceMessage.of(
|
IdUtil.fastSimpleUUID(),
|
IotDeviceMessageMethodEnum.EVENT_POST.getMethod(),
|
IotDeviceEventPostReqDTO.of(
|
"alarm",
|
MapUtil.<String, Object>builder()
|
.put("level", "warning")
|
.put("message", "temperature too high")
|
.put("threshold", 40)
|
.put("current", 42)
|
.build(),
|
System.currentTimeMillis()),
|
null, null, null);
|
// 2.2 序列化
|
byte[] payload = SERIALIZER.serialize(request);
|
String jsonMessage = StrUtil.utf8Str(payload);
|
log.info("[testEventPost][Serialize: {}, 请求消息: {}]", SERIALIZER.getType(), request);
|
|
// 3.1 发送并等待响应
|
String response = sendAndReceive(ws, jsonMessage);
|
// 3.2 解码响应
|
if (response != null) {
|
IotDeviceMessage responseMessage = SERIALIZER.deserialize(StrUtil.utf8Bytes(response));
|
log.info("[testEventPost][响应消息: {}]", responseMessage);
|
} else {
|
log.warn("[testEventPost][未收到响应]");
|
}
|
|
// 4. 关闭连接
|
ws.close();
|
}
|
|
// ===================== 辅助方法 =====================
|
|
/**
|
* 创建 WebSocket 连接(同步)
|
*
|
* @return WebSocket 连接
|
*/
|
private WebSocket createWebSocketConnection() throws Exception {
|
WebSocketClient wsClient = vertx.createWebSocketClient();
|
WebSocketConnectOptions options = new WebSocketConnectOptions()
|
.setHost(SERVER_HOST)
|
.setPort(SERVER_PORT)
|
.setURI(WS_PATH);
|
return wsClient.connect(options).toCompletionStage().toCompletableFuture().get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
|
}
|
|
/**
|
* 发送消息并等待响应(同步)
|
*
|
* @param ws WebSocket 连接
|
* @param message 请求消息
|
* @return 响应消息
|
*/
|
private String sendAndReceive(WebSocket ws, String message) throws Exception {
|
CountDownLatch latch = new CountDownLatch(1);
|
AtomicReference<String> responseRef = new AtomicReference<>();
|
|
// 设置消息处理器
|
ws.textMessageHandler(response -> {
|
log.info("[sendAndReceive][收到响应: {}]", response);
|
responseRef.set(response);
|
latch.countDown();
|
});
|
|
// 发送请求
|
log.info("[sendAndReceive][发送请求: {}]", message);
|
ws.writeTextMessage(message);
|
|
// 等待响应
|
boolean completed = latch.await(TIMEOUT_SECONDS, TimeUnit.SECONDS);
|
if (!completed) {
|
log.warn("[sendAndReceive][等待响应超时]");
|
}
|
return responseRef.get();
|
}
|
|
/**
|
* 执行子设备认证(同步)
|
*
|
* @param ws WebSocket 连接
|
* @return 认证响应消息
|
*/
|
private IotDeviceMessage authenticate(WebSocket ws) throws Exception {
|
IotDeviceAuthReqDTO authInfo = IotDeviceAuthUtils.getAuthInfo(PRODUCT_KEY, DEVICE_NAME, DEVICE_SECRET);
|
IotDeviceAuthReqDTO authReqDTO = new IotDeviceAuthReqDTO()
|
.setClientId(authInfo.getClientId())
|
.setUsername(authInfo.getUsername())
|
.setPassword(authInfo.getPassword());
|
IotDeviceMessage request = IotDeviceMessage.of(IdUtil.fastSimpleUUID(), "auth", authReqDTO, null, null, null);
|
|
byte[] payload = SERIALIZER.serialize(request);
|
String jsonMessage = StrUtil.utf8Str(payload);
|
log.info("[authenticate][发送认证请求: {}]", jsonMessage);
|
|
String response = sendAndReceive(ws, jsonMessage);
|
if (response != null) {
|
return SERIALIZER.deserialize(StrUtil.utf8Bytes(response));
|
}
|
return null;
|
}
|
|
}
|