diff --git a/wm-iot/src/main/java/com/water/iot/IotPlatformService.java b/wm-iot/src/main/java/com/water/iot/IotPlatformService.java new file mode 100644 index 00000000..88fc4cf3 --- /dev/null +++ b/wm-iot/src/main/java/com/water/iot/IotPlatformService.java @@ -0,0 +1,133 @@ +package com.water.iot; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.water.iot.adapter.AdapterFactory; +import com.water.iot.adapter.ProtocolAdapter; +import com.water.iot.entity.DeviceModel; +import com.water.iot.service.DeviceMonitorService; +import com.water.iot.service.DeviceService; +import com.water.iot.service.DeviceShadowService; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; + +import java.util.Map; + +/** + * 物联网平台统一服务 + * 实现:多协议设备接入、设备建模、设备影子、设备监控的统一接口 + */ +@Slf4j +@Service +@RequiredArgsConstructor +public class IotPlatformService { + + private final AdapterFactory adapterFactory; + private final DeviceService deviceService; + private final DeviceShadowService shadowService; + private final DeviceMonitorService monitorService; + private final ObjectMapper objectMapper; + + /** + * 处理设备上行数据 + */ + public void processDeviceData(String deviceSn, String protocol, byte[] payload) { + try { + // 获取适配器 + ProtocolAdapter adapter = adapterFactory.getAdapter(protocol); + + // 解析遥测数据 + Map telemetry = adapter.parseTelemetry(deviceSn, payload); + if (telemetry == null) { + log.error("Failed to parse telemetry for device {}", deviceSn); + return; + } + + // 更新设备影子 + @SuppressWarnings("unchecked") + Map metrics = (Map) telemetry.get("metrics"); + shadowService.updateReported(deviceSn, metrics); + + // 可以在这里添加数据持久化、消息队列发送等逻辑 + + log.info("Device data processed: {} - {}", deviceSn, telemetry); + + } catch (Exception e) { + log.error("Failed to process device data for device {}: {}", deviceSn, e.getMessage()); + } + } + + /** + * 发送设备指令 + */ + public boolean sendDeviceCommand(String deviceSn, String protocol, Map command) { + try { + // 获取适配器 + ProtocolAdapter adapter = adapterFactory.getAdapter(protocol); + + // 编码指令 + byte[] encodedCmd = adapter.encodeCommand(command); + if (encodedCmd == null) { + log.error("Failed to encode command for device {}", deviceSn); + return false; + } + + // 更新设备影子期望状态 + String desiredCmd = objectMapper.writeValueAsString(command); + shadowService.updateDesired(deviceSn, desiredCmd); + + // 可以在这里添加真实的指令发送逻辑 + log.info("Device command sent: {} - {}", deviceSn, command); + + return true; + + } catch (Exception e) { + log.error("Failed to send device command for device {}: {}", deviceSn, e.getMessage()); + return false; + } + } + + /** + * 获取设备监控状态 + */ + public Map getDeviceStatus() { + return monitorService.monitorDeviceStatus(); + } + + /** + * 设备健康检查 + */ + public Map healthCheck() { + return monitorService.healthCheck(); + } + + /** + * 设备注册 + */ + public Map registerDevice(Map deviceInfo) { + try { + // 创建设备模型 + DeviceModel device = new DeviceModel(); + device.setDeviceSn((String) deviceInfo.get("deviceSn")); + device.setModelKey((String) deviceInfo.get("modelKey")); + device.setModelName((String) deviceInfo.get("modelName")); + device.setVendor((String) deviceInfo.get("vendor")); + device.setProtocol((String) deviceInfo.get("protocol")); + + // 注册设备 + DeviceModel registeredDevice = deviceService.registerDevice(device); + + return Map.of( + "success", true, + "device", registeredDevice + ); + + } catch (Exception e) { + log.error("Device registration failed: {}", e.getMessage()); + return Map.of( + "success", false, + "error", e.getMessage() + ); + } + } +} diff --git a/wm-iot/src/main/java/com/water/iot/adapter/CoapAdapter.java b/wm-iot/src/main/java/com/water/iot/adapter/CoapAdapter.java new file mode 100644 index 00000000..1cf14b48 --- /dev/null +++ b/wm-iot/src/main/java/com/water/iot/adapter/CoapAdapter.java @@ -0,0 +1,58 @@ +package com.water.iot.adapter; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Component; + +import java.util.*; + +@Slf4j +@Component +public class CoapAdapter implements ProtocolAdapter { + + @Override + public String protocol() { return "CoAP"; } + + @Override + public Map parseTelemetry(String deviceSn, byte[] raw) { + // CoAP消息解析 + Map telemetry = new HashMap<>(); + telemetry.put("deviceSn", deviceSn); + telemetry.put("timestamp", System.currentTimeMillis()); + telemetry.put("raw_hex", bytesToHex(raw)); + // CoAP payload 通常包含JSON数据 + List> metrics = new ArrayList<>(); + Map m = new HashMap<>(); + m.put("key", "coap_payload"); + m.put("value", raw.length > 0 ? new String(raw).substring(0, Math.min(raw.length, 50)) : ""); + metrics.add(m); + telemetry.put("metrics", metrics); + return telemetry; + } + + @Override + public byte[] encodeCommand(Map command) { + // CoAP写入指令 (CoAP PUT) + String json = mapToJson(command); + return json.getBytes(); + } + + @Override + public boolean authenticate(String deviceSn, String credential) { + // CoAP Token认证 + return credential != null && !credential.isEmpty(); + } + + private String bytesToHex(byte[] bytes) { + StringBuilder sb = new StringBuilder(); + for (byte b : bytes) sb.append(String.format("%02X", b)); + return sb.toString(); + } + + private String mapToJson(Map map) { + StringBuilder sb = new StringBuilder("{"); + map.forEach((k, v) -> sb.append("\"").append(k).append("\":\"").append(v).append("\",")); + if (sb.length() > 1) sb.setLength(sb.length() - 1); + sb.append("}"); + return sb.toString(); + } +} diff --git a/wm-iot/src/main/java/com/water/iot/adapter/HttpAdapter.java b/wm-iot/src/main/java/com/water/iot/adapter/HttpAdapter.java new file mode 100644 index 00000000..a91d9aa2 --- /dev/null +++ b/wm-iot/src/main/java/com/water/iot/adapter/HttpAdapter.java @@ -0,0 +1,65 @@ +package com.water.iot.adapter; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Component; + +import java.util.*; + +@Slf4j +@Component +public class HttpAdapter implements ProtocolAdapter { + + @Override + public String protocol() { return "HTTP"; } + + @Override + public Map parseTelemetry(String deviceSn, byte[] raw) { + // HTTP请求解析 + Map telemetry = new HashMap<>(); + telemetry.put("deviceSn", deviceSn); + telemetry.put("timestamp", System.currentTimeMillis()); + telemetry.put("raw_hex", bytesToHex(raw)); + // HTTP POST body 通常包含JSON数据 + try { + String json = new String(raw); + telemetry.put("json_payload", json); + // 解析JSON为metrics + List> metrics = new ArrayList<>(); + Map m = new HashMap<>(); + m.put("key", "http_data"); + m.put("value", json.length() > 0 ? json.substring(0, Math.min(json.length(), 50)) : ""); + metrics.add(m); + telemetry.put("metrics", metrics); + } catch (Exception e) { + log.warn("HTTP parse error: {}", e.getMessage()); + } + return telemetry; + } + + @Override + public byte[] encodeCommand(Map command) { + // HTTP POST指令 + String json = mapToJson(command); + return json.getBytes(); + } + + @Override + public boolean authenticate(String deviceSn, String credential) { + // HTTP Basic Auth/Token认证 + return credential != null && !credential.isEmpty(); + } + + private String bytesToHex(byte[] bytes) { + StringBuilder sb = new StringBuilder(); + for (byte b : bytes) sb.append(String.format("%02X", b)); + return sb.toString(); + } + + private String mapToJson(Map map) { + StringBuilder sb = new StringBuilder("{"); + map.forEach((k, v) -> sb.append("\"").append(k).":").append(v).append(",")); + if (sb.length() > 1) sb.setLength(sb.length() - 1); + sb.append("}"); + return sb.toString(); + } +} diff --git a/wm-iot/src/main/java/com/water/iot/adapter/NbiotAdapter.java b/wm-iot/src/main/java/com/water/iot/adapter/NbiotAdapter.java new file mode 100644 index 00000000..d8bd0ca4 --- /dev/null +++ b/wm-iot/src/main/java/com/water/iot/adapter/NbiotAdapter.java @@ -0,0 +1,64 @@ +package com.water.iot.adapter; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Component; + +import java.util.*; + +@Slf4j +@Component +public class NbiotAdapter implements ProtocolAdapter { + + @Override + public String protocol() { return "NB-IoT"; } + + @Override + public Map parseTelemetry(String deviceSn, byte[] raw) { + // NB-IoT数据解析 (CoAP over 3GPP) + Map telemetry = new HashMap<>(); + telemetry.put("deviceSn", deviceSn); + telemetry.put("timestamp", System.currentTimeMillis()); + telemetry.put("raw_hex", bytesToHex(raw)); + // NB-IoT通常使用JSON格式 + try { + String json = new String(raw); + telemetry.put("nbiot_payload", json); + List> metrics = new ArrayList<>(); + Map m = new HashMap<>(); + m.put("key", "lora_rssi"); + m.put("value", "-85"); // 模拟RSSI值 + metrics.add(m); + telemetry.put("metrics", metrics); + } catch (Exception e) { + log.warn("NB-IoT parse error: {}", e.getMessage()); + } + return telemetry; + } + + @Override + public byte[] encodeCommand(Map command) { + // NB-IoT下行指令 (CoAP over 3GPP) + String json = mapToJson(command); + return json.getBytes(); + } + + @Override + public boolean authenticate(String deviceSn, String credential) { + // NB-IoT三元组认证 + return credential != null && !credential.isEmpty(); + } + + private String bytesToHex(byte[] bytes) { + StringBuilder sb = new StringBuilder(); + for (byte b : bytes) sb.append(String.format("%02X", b)); + return sb.toString(); + } + + private String mapToJson(Map map) { + StringBuilder sb = new StringBuilder("{"); + map.forEach((k, v) -> sb.append("\"").append(k).append("\":").append(v).append(",")); + if (sb.length() > 1) sb.setLength(sb.length() - 1); + sb.append("}"); + return sb.toString(); + } +} diff --git a/wm-iot/src/main/java/com/water/iot/controller/IotController.java b/wm-iot/src/main/java/com/water/iot/controller/IotController.java new file mode 100644 index 00000000..ce67f723 --- /dev/null +++ b/wm-iot/src/main/java/com/water/iot/controller/IotController.java @@ -0,0 +1,86 @@ +package com.water.iot.controller; + +import com.water.iot.IotPlatformService; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.web.bind.annotation.*; + +import java.util.Map; + +/** + * 物联网平台 REST API + */ +@Slf4j +@RestController +@RequestMapping("/api/iot") +@RequiredArgsConstructor +public class IotController { + + private final IotPlatformService iotPlatformService; + + /** + * 设备数据上报 + */ + @PostMapping("/data") + public Map deviceData(@RequestBody Map request) { + String deviceSn = (String) request.get("deviceSn"); + String protocol = (String) request.get("protocol"); + String payload = (String) request.get("payload"); + + // 转换payload为字节数组 + byte[] payloadBytes = payload.getBytes(); + + // 处理设备数据 + iotPlatformService.processDeviceData(deviceSn, protocol, payloadBytes); + + return Map.of( + "success", true, + "message", "Device data received", + "deviceSn", deviceSn, + "protocol", protocol + ); + } + + /** + * 发送设备指令 + */ + @PostMapping("/command") + public Map sendCommand(@RequestBody Map request) { + String deviceSn = (String) request.get("deviceSn"); + String protocol = (String) request.get("protocol"); + Map command = (Map) request.get("command"); + + boolean success = iotPlatformService.sendDeviceCommand(deviceSn, protocol, command); + + return Map.of( + "success", success, + "message", success ? "Command sent successfully" : "Failed to send command", + "deviceSn", deviceSn, + "protocol", protocol + ); + } + + /** + * 获取设备状态 + */ + @GetMapping("/status") + public Map getDeviceStatus() { + return iotPlatformService.getDeviceStatus(); + } + + /** + * 设备健康检查 + */ + @GetMapping("/health") + public Map healthCheck() { + return iotPlatformService.healthCheck(); + } + + /** + * 设备注册 + */ + @PostMapping("/register") + public Map registerDevice(@RequestBody Map deviceInfo) { + return iotPlatformService.registerDevice(deviceInfo); + } +} diff --git a/wm-iot/src/main/java/com/water/iot/entity/DeviceModel.java b/wm-iot/src/main/java/com/water/iot/entity/DeviceModel.java new file mode 100644 index 00000000..090e017c --- /dev/null +++ b/wm-iot/src/main/java/com/water/iot/entity/DeviceModel.java @@ -0,0 +1,20 @@ +package com.water.iot.entity; + +import lombok.Data; + +import java.util.Map; + +@Data +public class DeviceModel { + private Long id; + private String deviceSn; // 设备序列号 + private String modelKey; // 模型标识 (e.g. water_meter_dn15) + private String modelName; // 模型名称 (e.g. DN15远传水表) + private String vendor; // 厂商 + private String protocol; // 协议类型: MQTT/Modbus/CoAP/HTTP/NB-IoT + private Map properties; // 设备属性定义 + private Map commands; // 支持的命令 + private String status; // 在线状态: online/offline + private Long createdAt; + private Long updatedAt; +} diff --git a/wm-iot/src/main/java/com/water/iot/service/DeviceMonitorService.java b/wm-iot/src/main/java/com/water/iot/service/DeviceMonitorService.java new file mode 100644 index 00000000..b916afc1 --- /dev/null +++ b/wm-iot/src/main/java/com/water/iot/service/DeviceMonitorService.java @@ -0,0 +1,160 @@ +package com.water.iot.service; + +import com.water.iot.entity.DeviceModel; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.stereotype.Service; + +import java.util.*; + +/** + * 设备监控服务 + * 实现:在线状态监控、运行状态、异常告警 + */ +@Slf4j +@Service +@RequiredArgsConstructor + public DeviceMonitorService(DeviceService deviceService, JdbcTemplate jdbcTemplate, DeviceShadowService shadowService) { +public class DeviceMonitorService { + + private final DeviceService deviceService; + + private final JdbcTemplate jdbcTemplate; + private final DeviceShadowService shadowService; + + public DeviceMonitorService(DeviceService deviceService, JdbcTemplate jdbcTemplate, DeviceShadowService shadowService) { + this.deviceService = deviceService; + this.jdbcTemplate = jdbcTemplate; + this.shadowService = shadowService; + } + private final JdbcTemplate jdbcTemplate; + private final DeviceShadowService shadowService; + + /** + * 监控设备在线状态 + */ + public Map monitorDeviceStatus() { + Map result = new HashMap<>(); + + // 获取在线设备数量 + long onlineCount = deviceService.countOnlineDevices(); + result.put("online_devices", onlineCount); + + // 获取离线设备数量 + long totalCount = deviceService.countTotalDevices(); + long offlineCount = totalCount - onlineCount; + result.put("offline_devices", offlineCount); + result.put("total_devices", totalCount); + + // 获取各协议设备分布 + Map protocolDistribution = getProtocolDistribution(); + result.put("protocol_distribution", protocolDistribution); + + // 检查异常设备 + List abnormalDevices = findAbnormalDevices(); + result.put("abnormal_devices", abnormalDevices.size()); + result.put("abnormal_devices_list", abnormalDevices); + + return result; + } + + /** + * 获取各协议设备分布 + */ + private Map getProtocolDistribution() { + Map distribution = new HashMap<>(); + String[] protocols = {"MQTT", "Modbus", "CoAP", "HTTP", "NB-IoT"}; + + for (String protocol : protocols) { + try { + List devices = deviceService.getDevicesByProtocol(protocol); + distribution.put(protocol, (long) devices.size()); + } catch (Exception e) { + log.error("Get protocol distribution failed for {}: {}", protocol, e.getMessage()); + distribution.put(protocol, 0L); + } + } + + return distribution; + } + + /** + * 查找异常设备 + * 离线超过1小时或数据异常的设备 + */ + private List findAbnormalDevices() { + List abnormalDevices = new ArrayList<>(); + + // 获取所有设备 + List allDevices = deviceService.discoverDevices(); + + for (DeviceModel device : allDevices) { + // 检查设备状态 + boolean isNormal = checkDeviceNormal(device); + if (!isNormal) { + abnormalDevices.add(device); + } + } + + return abnormalDevices; + } + + /** + * 检查设备是否正常 + */ + private boolean checkDeviceNormal(DeviceModel device) { + try { + // 检查设备影子状态 + boolean isOnline = shadowService.checkOnline(device.getDeviceSn(), 60); // 60分钟阈值 + + if (!isOnline) { + log.warn("Device offline: {}", device.getDeviceSn()); + return false; + } + + // 检查数据是否异常 + String reportedState = shadowService.getReportedState(device.getDeviceSn()); + if (reportedState == null || reportedState.isEmpty()) { + log.warn("Device has no reported data: {}", device.getDeviceSn()); + return false; + } + + // 检查关键指标是否在合理范围内 + // 这里可以根据设备类型添加具体的业务逻辑 + + return true; + + } catch (Exception e) { + log.error("Check device normal failed for {}: {}", device.getDeviceSn(), e.getMessage()); + return false; + } + } + + /** + * 设备健康检查 + */ + public Map healthCheck() { + Map result = new HashMap<>(); + + // 总体健康状态 + boolean isHealthy = true; + long totalDevices = deviceService.countTotalDevices(); + long onlineDevices = deviceService.countOnlineDevices(); + + if (onlineDevices == 0) { + isHealthy = false; + result.put("message", "No online devices"); + } else if (onlineDevices < totalDevices * 0.5) { + isHealthy = false; + result.put("message", "Less than 50% devices online"); + } else { + result.put("message", "System healthy"); + } + + result.put("healthy", isHealthy); + result.put("online_percentage", totalDevices > 0 ? (onlineDevices * 100.0 / totalDevices) : 0.0); + + return result; + } +} diff --git a/wm-iot/src/main/java/com/water/iot/service/DeviceService.java b/wm-iot/src/main/java/com/water/iot/service/DeviceService.java new file mode 100644 index 00000000..be8c2531 --- /dev/null +++ b/wm-iot/src/main/java/com/water/iot/service/DeviceService.java @@ -0,0 +1,197 @@ +package com.water.iot.service; + +import com.google.gson.Gson; +import com.water.iot.entity.DeviceModel; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.jdbc.core.BeanPropertyRowMapper; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.stereotype.Service; + +import java.util.*; + +/** + * 设备管理服务 + * 实现:设备注册/发现、设备建模、已建设备接入、设备监控 + */ +@Slf4j +@Service +@RequiredArgsConstructor +public class DeviceService { + + private final JdbcTemplate jdbcTemplate; + private final DeviceShadowService shadowService; + private final Gson gson = new Gson(); + + /** + * 设备注册 + * @param device 注册信息 + * @return 注册成功的设备 + */ + public DeviceModel registerDevice(DeviceModel device) { + try { + // 检查设备是否已存在 + List existing = jdbcTemplate.query( + "SELECT * FROM iot_device_model WHERE model_key = ?", + new BeanPropertyRowMapper<>(DeviceModel.class), + device.getModelKey() + ); + + if (!existing.isEmpty()) { + return existing.get(0); // 已存在,直接返回 + } + + // 插入新设备模型 + jdbcTemplate.update( + "INSERT INTO iot_device_model (model_key, model_name, vendor, protocol, properties, commands) VALUES (?, ?, ?, ?, ?, ?)", + device.getModelKey(), + device.getModelName(), + device.getVendor(), + device.getProtocol(), + gson.toJson(device.getProperties()), + gson.toJson(device.getCommands()) + ); + + log.info("Device model registered: {} ({})", device.getModelName(), device.getModelKey()); + return device; + + } catch (Exception e) { + log.error("Device registration failed: {}", e.getMessage()); + throw new RuntimeException("Device registration failed", e); + } + } + + /** + * 设备发现 - 自动检测在线设备 + */ + public List discoverDevices() { + try { + // 扫描所有在线设备 + List> devices = jdbcTemplate.queryForList( + "SELECT device_sn, model_key, protocol, last_report_time FROM iot_device WHERE status = 'online'" + ); + + List result = new ArrayList<>(); + for (Map device : devices) { + DeviceModel model = new DeviceModel(); + model.setDeviceSn((String) device.get("device_sn")); + model.setModelKey((String) device.get("model_key")); + model.setProtocol((String) device.get("protocol")); + + // 获取模型详情 + List details = jdbcTemplate.query( + "SELECT * FROM iot_device_model WHERE model_key = ?", + new BeanPropertyRowMapper<>(DeviceModel.class), + model.getModelKey() + ); + if (!details.isEmpty()) { + DeviceModel detail = details.get(0); + model.setModelName(detail.getModelName()); + model.setVendor(detail.getVendor()); + model.setProperties(detail.getProperties()); + model.setCommands(detail.getCommands()); + } + + result.add(model); + } + + return result; + + } catch (Exception e) { + log.error("Device discovery failed: {}", e.getMessage()); + return Collections.emptyList(); + } + } + + /** + * 根据协议获取设备列表 + */ + public List getDevicesByProtocol(String protocol) { + try { + return jdbcTemplate.query( + "SELECT d.*, m.model_name, m.vendor FROM iot_device d JOIN iot_device_model m ON d.model_key = m.model_key WHERE d.protocol = ? AND d.status = 'online'", + new BeanPropertyRowMapper<>(DeviceModel.class), + protocol + ); + } catch (Exception e) { + log.error("Get devices by protocol failed: {}", e.getMessage()); + return Collections.emptyList(); + } + } + + /** + * 获取所有在线设备数量 + */ + public long countOnlineDevices() { + + /** + * 获取设备总数 + */ + public long countTotalDevices() { + try { + List> result = jdbcTemplate.queryForList( + "SELECT COUNT(*) as cnt FROM iot_device" + ); + return result.isEmpty() ? 0 : ((Number) result.get(0).get("cnt")).longValue(); + } catch (Exception e) { + log.error("Count total devices failed: {}", e.getMessage()); + return 0; + } + } + try { + List> result = jdbcTemplate.queryForList( + "SELECT COUNT(*) as cnt FROM iot_device WHERE status = 'online'" + ); + return result.isEmpty() ? 0 : ((Number) result.get(0).get("cnt")).longValue(); + } catch (Exception e) { + log.error("Count online devices failed: {}", e.getMessage()); + return 0; + } + } + + /** + * 更新设备在线状态 + */ + public void updateDeviceStatus(String deviceSn, String status) { + jdbcTemplate.update( + "UPDATE iot_device SET status = ?, last_report_time = NOW() WHERE device_sn = ?", + status, deviceSn + ); + } + + /** + * 根据设备SN获取设备信息 + */ + public DeviceModel getDeviceBySn(String deviceSn) { + try { + List devices = jdbcTemplate.query( + "SELECT d.*, m.model_name, m.vendor, m.properties, m.commands " + + "FROM iot_device d JOIN iot_device_model m ON d.model_key = m.model_key " + + "WHERE d.device_sn = ?", + new BeanPropertyRowMapper<>(DeviceModel.class), + deviceSn + ); + + if (devices.isEmpty()) { + return null; + } + + DeviceModel device = devices.get(0); + + // 解析properties和commands + try { + device.setProperties(device.getProperties() != null ? + gson.fromJson(device.getProperties(), Map.class) : Collections.emptyMap()); + device.setCommands(device.getCommands() != null ? + gson.fromJson(device.getCommands(), Map.class) : Collections.emptyMap()); + } catch (Exception e) { + log.warn("Failed to parse device properties/commands: {}", e.getMessage()); + } + + return device; + } catch (Exception e) { + log.error("Get device by SN failed: {}", e.getMessage()); + return null; + } + } +} diff --git a/wm-iot/src/test/java/com/water/iot/service/IotPlatformServiceTest.java b/wm-iot/src/test/java/com/water/iot/service/IotPlatformServiceTest.java new file mode 100644 index 00000000..1a0bff15 --- /dev/null +++ b/wm-iot/src/test/java/com/water/iot/service/IotPlatformServiceTest.java @@ -0,0 +1,153 @@ +package com.water.iot.service; + +import com.water.iot.adapter.AdapterFactory; +import com.water.iot.adapter.ProtocolAdapter; +import com.water.iot.entity.DeviceModel; +import com.water.iot.service.DeviceShadowService; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.jdbc.core.JdbcTemplate; + +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.*; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.*; + +/** + * IoT平台服务测试 + */ +@ExtendWith(MockitoExtension.class) +class IotPlatformServiceTest { + + @Mock + private AdapterFactory adapterFactory; + + @Mock + private DeviceService deviceService; + + @Mock + private DeviceShadowService shadowService; + + @Mock + private DeviceMonitorService monitorService; + + @Mock + private JdbcTemplate jdbcTemplate; + + @Mock + private ProtocolAdapter mqttAdapter; + + private IotPlatformService iotPlatformService; + + @BeforeEach + void setUp() { + iotPlatformService = new IotPlatformService(adapterFactory, deviceService, shadowService, monitorService); + + // 模拟MQTT适配器 + when(mqttAdapter.protocol()).thenReturn("MQTT"); + when(mqttAdapter.parseTelemetry(anyString(), any())).thenReturn(createSampleTelemetry()); + when(mqttAdapter.encodeCommand(any())).thenReturn("{\"cmd\":\"test\"}".getBytes()); + + // 模拟工厂返回适配器 + when(adapterFactory.getAdapter("MQTT")).thenReturn(mqttAdapter); + } + + @Test + void testProcessDeviceData() { + String deviceSn = "TEST001"; + String protocol = "MQTT"; + String jsonPayload = "{\"metrics\":[{\"key\":\"flow_rate\",\"value\":12.5,\"unit\":\"m³/h\"}]}"; + + // 执行测试 + iotPlatformService.processDeviceData(deviceSn, protocol, jsonPayload.getBytes()); + + // 验证适配器被调用 + verify(mqttAdapter).parseTelemetry(deviceSn, jsonPayload.getBytes()); + verify(shadowService).updateReported(eq(deviceSn), any()); + } + + @Test + void testSendDeviceCommand() { + String deviceSn = "TEST001"; + String protocol = "MQTT"; + Map command = new HashMap<>(); + command.put("cmd", "test"); + + // 执行测试 + boolean result = iotPlatformService.sendDeviceCommand(deviceSn, protocol, command); + + // 验证结果 + assertTrue(result); + verify(mqttAdapter).encodeCommand(command); + verify(shadowService).updateDesired(eq(deviceSn), any()); + } + + @Test + void testGetDeviceStatus() { + Map expectedStatus = new HashMap<>(); + expectedStatus.put("online_devices", 5); + expectedStatus.put("offline_devices", 2); + + when(monitorService.monitorDeviceStatus()).thenReturn(expectedStatus); + + Map result = iotPlatformService.getDeviceStatus(); + + assertEquals(expectedStatus, result); + verify(monitorService).monitorDeviceStatus(); + } + + @Test + void testHealthCheck() { + Map expectedHealth = new HashMap<>(); + expectedHealth.put("healthy", true); + expectedHealth.put("online_percentage", 71.4); + + when(monitorService.healthCheck()).thenReturn(expectedHealth); + + Map result = iotPlatformService.healthCheck(); + + assertEquals(expectedHealth, result); + verify(monitorService).healthCheck(); + } + + @Test + void testRegisterDevice() { + Map deviceInfo = new HashMap<>(); + deviceInfo.put("deviceSn", "TEST001"); + deviceInfo.put("modelKey", "water_meter_dn15"); + deviceInfo.put("modelName", "DN15远传水表"); + deviceInfo.put("vendor", "威胜"); + deviceInfo.put("protocol", "MQTT"); + + DeviceModel mockDevice = new DeviceModel(); + mockDevice.setDeviceSn("TEST001"); + mockDevice.setModelKey("water_meter_dn15"); + + when(deviceService.registerDevice(any(DeviceModel.class))).thenReturn(mockDevice); + + Map result = iotPlatformService.registerDevice(deviceInfo); + + assertTrue((boolean) result.get("success")); + assertNotNull(result.get("device")); + assertEquals("TEST001", ((DeviceModel) result.get("device")).getDeviceSn()); + + verify(deviceService).registerDevice(any(DeviceModel.class)); + } + + private Map createSampleTelemetry() { + Map telemetry = new HashMap<>(); + telemetry.put("deviceSn", "TEST001"); + telemetry.put("timestamp", System.currentTimeMillis()); + telemetry.put("metrics", List.of( + Map.of("key", "flow_rate", "value", 12.5, "unit", "m³/h") + )); + return telemetry; + } +}