farm模块-mqtt整合及优化
This commit is contained in:
@@ -1,38 +1,35 @@
|
||||
package cn.iocoder.yudao.module.farm.controller.admin.deviceinfo;
|
||||
|
||||
import cn.iocoder.yudao.framework.apilog.core.annotation.ApiAccessLog;
|
||||
import cn.iocoder.yudao.framework.common.pojo.CommonResult;
|
||||
import cn.iocoder.yudao.framework.common.pojo.PageParam;
|
||||
import cn.iocoder.yudao.framework.common.pojo.PageResult;
|
||||
import cn.iocoder.yudao.framework.common.util.object.BeanUtils;
|
||||
import cn.iocoder.yudao.framework.excel.core.util.ExcelUtils;
|
||||
import cn.iocoder.yudao.module.farm.controller.admin.deviceinfo.vo.DeviceInfoPageReqVO;
|
||||
import cn.iocoder.yudao.module.farm.controller.admin.deviceinfo.vo.DeviceInfoRespVO;
|
||||
import cn.iocoder.yudao.module.farm.controller.admin.deviceinfo.vo.DeviceInfoSaveReqVO;
|
||||
import org.springframework.web.bind.annotation.*;
|
||||
import jakarta.validation.Valid;
|
||||
import jakarta.validation.Valid;
|
||||
import jakarta.validation.Valid;
|
||||
import jakarta.servlet.http.HttpServletResponse;
|
||||
import java.util.List;
|
||||
import static cn.iocoder.yudao.framework.apilog.core.enums.OperateTypeEnum.EXPORT;
|
||||
|
||||
import jakarta.annotation.Resource;
|
||||
import org.springframework.validation.annotation.Validated;
|
||||
|
||||
import io.swagger.v3.oas.annotations.tags.Tag;
|
||||
import io.swagger.v3.oas.annotations.Parameter;
|
||||
import cn.iocoder.yudao.module.farm.dal.dataobject.deviceinfo.DeviceInfoDO;
|
||||
import cn.iocoder.yudao.module.farm.dal.dataobject.mqtt.MqttConnectionDO;
|
||||
import cn.iocoder.yudao.module.farm.service.deviceinfo.DeviceInfoService;
|
||||
import cn.iocoder.yudao.module.farm.service.mqtt.DeviceMqttService;
|
||||
import io.swagger.v3.oas.annotations.Operation;
|
||||
import io.swagger.v3.oas.annotations.Parameter;
|
||||
import io.swagger.v3.oas.annotations.tags.Tag;
|
||||
import jakarta.annotation.Resource;
|
||||
import jakarta.servlet.http.HttpServletResponse;
|
||||
import jakarta.validation.Valid;
|
||||
import org.eclipse.paho.client.mqttv3.MqttException;
|
||||
import org.springframework.validation.annotation.Validated;
|
||||
import org.springframework.web.bind.annotation.*;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.List;
|
||||
|
||||
import cn.iocoder.yudao.framework.common.pojo.PageParam;
|
||||
import cn.iocoder.yudao.framework.common.pojo.PageResult;
|
||||
import cn.iocoder.yudao.framework.common.pojo.CommonResult;
|
||||
import cn.iocoder.yudao.framework.common.util.object.BeanUtils;
|
||||
import static cn.iocoder.yudao.framework.apilog.core.enums.OperateTypeEnum.EXPORT;
|
||||
import static cn.iocoder.yudao.framework.common.pojo.CommonResult.error;
|
||||
import static cn.iocoder.yudao.framework.common.pojo.CommonResult.success;
|
||||
|
||||
import cn.iocoder.yudao.framework.excel.core.util.ExcelUtils;
|
||||
|
||||
import cn.iocoder.yudao.framework.apilog.core.annotation.ApiAccessLog;
|
||||
|
||||
import cn.iocoder.yudao.module.farm.dal.dataobject.deviceinfo.DeviceInfoDO;
|
||||
import cn.iocoder.yudao.module.farm.service.deviceinfo.DeviceInfoService;
|
||||
|
||||
@Tag(name = "管理后台 - 设备信息")
|
||||
@RestController
|
||||
@RequestMapping("/farm/device-info")
|
||||
@@ -41,6 +38,8 @@ public class DeviceInfoController {
|
||||
|
||||
@Resource
|
||||
private DeviceInfoService deviceInfoService;
|
||||
@Resource
|
||||
private DeviceMqttService deviceMqttService;
|
||||
|
||||
@PostMapping("/create")
|
||||
@Operation(summary = "创建设备信息")
|
||||
@@ -98,4 +97,15 @@ public class DeviceInfoController {
|
||||
BeanUtils.toBean(list, DeviceInfoRespVO.class));
|
||||
}
|
||||
|
||||
@PostMapping("/startSub")
|
||||
@Operation(summary = "订阅设备上报并持久化")
|
||||
public CommonResult<Boolean> startSub(@Valid @RequestBody MqttConnectionDO createReqVO) {
|
||||
try {
|
||||
deviceMqttService.subscribeTelemetryAndPersist(createReqVO.getDeviceId());
|
||||
return success(true);
|
||||
} catch (MqttException e) {
|
||||
return error(500, "设备连接失败");
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,204 @@
|
||||
package cn.iocoder.yudao.module.farm.controller.admin.devicestatus;
|
||||
|
||||
import cn.iocoder.yudao.framework.common.pojo.CommonResult;
|
||||
import cn.iocoder.yudao.framework.common.util.object.BeanUtils;
|
||||
import cn.iocoder.yudao.module.farm.controller.admin.devicestatus.vo.DeviceStatusRespVO;
|
||||
import cn.iocoder.yudao.module.farm.controller.admin.devicestatus.vo.DeviceStatusStatisticsRespVO;
|
||||
import cn.iocoder.yudao.module.farm.dal.dataobject.deviceinfo.DeviceInfoDO;
|
||||
import cn.iocoder.yudao.module.farm.dal.mysql.deviceinfo.DeviceInfoMapper;
|
||||
import cn.iocoder.yudao.module.farm.enums.DeviceStateEnum;
|
||||
import cn.iocoder.yudao.module.farm.service.devicestatus.DeviceStatusService;
|
||||
import io.swagger.v3.oas.annotations.Operation;
|
||||
import io.swagger.v3.oas.annotations.Parameter;
|
||||
import io.swagger.v3.oas.annotations.tags.Tag;
|
||||
import jakarta.annotation.Resource;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.security.access.prepost.PreAuthorize;
|
||||
import org.springframework.validation.annotation.Validated;
|
||||
import org.springframework.web.bind.annotation.*;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
import static cn.iocoder.yudao.framework.common.pojo.CommonResult.success;
|
||||
|
||||
@Tag(name = "管理后台 - 设备状态管理")
|
||||
@RestController
|
||||
@RequestMapping("/farm/device-status")
|
||||
@Validated
|
||||
@Slf4j
|
||||
public class DeviceStatusController {
|
||||
|
||||
@Resource
|
||||
private DeviceStatusService deviceStatusService;
|
||||
|
||||
@Resource
|
||||
private DeviceInfoMapper deviceInfoMapper;
|
||||
|
||||
@GetMapping("/check/{deviceUniqueId}")
|
||||
@Operation(summary = "检查设备是否在线")
|
||||
@Parameter(name = "deviceUniqueId", description = "设备唯一标识", required = true, example = "2409AA01")
|
||||
@PreAuthorize("@ss.hasPermission('farm:device-status:query')")
|
||||
public CommonResult<Boolean> checkDeviceOnline(@PathVariable("deviceUniqueId") String deviceUniqueId) {
|
||||
boolean online = deviceStatusService.isDeviceOnline(deviceUniqueId);
|
||||
return success(online);
|
||||
}
|
||||
|
||||
@GetMapping("/state/{deviceUniqueId}")
|
||||
@Operation(summary = "获取设备状态")
|
||||
@Parameter(name = "deviceUniqueId", description = "设备唯一标识", required = true, example = "2409AA01")
|
||||
@PreAuthorize("@ss.hasPermission('farm:device-status:query')")
|
||||
public CommonResult<DeviceStatusRespVO> getDeviceStatus(@PathVariable("deviceUniqueId") String deviceUniqueId) {
|
||||
DeviceInfoDO device = deviceInfoMapper.selectOne(DeviceInfoDO::getDeviceUniqueId, deviceUniqueId);
|
||||
if (device == null) {
|
||||
return success(null);
|
||||
}
|
||||
|
||||
DeviceStatusRespVO respVO = BeanUtils.toBean(device, DeviceStatusRespVO.class);
|
||||
if (device.getState() != null) {
|
||||
DeviceStateEnum stateEnum = DeviceStateEnum.valueOf(device.getState());
|
||||
if (stateEnum != null) {
|
||||
respVO.setStateName(stateEnum.getName());
|
||||
respVO.setOnline(DeviceStateEnum.isOnline(device.getState()));
|
||||
}
|
||||
}
|
||||
return success(respVO);
|
||||
}
|
||||
|
||||
@GetMapping("/online-devices")
|
||||
@Operation(summary = "获取在线设备列表")
|
||||
@PreAuthorize("@ss.hasPermission('farm:device-status:query')")
|
||||
public CommonResult<List<DeviceStatusRespVO>> getOnlineDevices() {
|
||||
List<DeviceInfoDO> devices = deviceStatusService.getOnlineDevices();
|
||||
List<DeviceStatusRespVO> respList = BeanUtils.toBean(devices, DeviceStatusRespVO.class);
|
||||
|
||||
// 填充状态信息
|
||||
respList.forEach(resp -> {
|
||||
if (resp.getState() != null) {
|
||||
DeviceStateEnum stateEnum = DeviceStateEnum.valueOf(resp.getState());
|
||||
if (stateEnum != null) {
|
||||
resp.setStateName(stateEnum.getName());
|
||||
resp.setOnline(DeviceStateEnum.isOnline(resp.getState()));
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
return success(respList);
|
||||
}
|
||||
|
||||
@GetMapping("/offline-devices")
|
||||
@Operation(summary = "获取离线设备列表")
|
||||
@PreAuthorize("@ss.hasPermission('farm:device-status:query')")
|
||||
public CommonResult<List<DeviceStatusRespVO>> getOfflineDevices() {
|
||||
List<DeviceInfoDO> devices = deviceStatusService.getOfflineDevices();
|
||||
List<DeviceStatusRespVO> respList = BeanUtils.toBean(devices, DeviceStatusRespVO.class);
|
||||
|
||||
// 填充状态信息
|
||||
respList.forEach(resp -> {
|
||||
if (resp.getState() != null) {
|
||||
DeviceStateEnum stateEnum = DeviceStateEnum.valueOf(resp.getState());
|
||||
if (stateEnum != null) {
|
||||
resp.setStateName(stateEnum.getName());
|
||||
resp.setOnline(DeviceStateEnum.isOnline(resp.getState()));
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
return success(respList);
|
||||
}
|
||||
|
||||
@GetMapping("/timeout-devices")
|
||||
@Operation(summary = "获取超时未上报的设备列表")
|
||||
@Parameter(name = "timeoutMinutes", description = "超时分钟数", required = false, example = "10")
|
||||
@PreAuthorize("@ss.hasPermission('farm:device-status:query')")
|
||||
public CommonResult<List<DeviceStatusRespVO>> getTimeoutDevices(
|
||||
@RequestParam(value = "timeoutMinutes", defaultValue = "10") int timeoutMinutes) {
|
||||
List<DeviceInfoDO> devices = deviceStatusService.getTimeoutDevices(timeoutMinutes);
|
||||
List<DeviceStatusRespVO> respList = BeanUtils.toBean(devices, DeviceStatusRespVO.class);
|
||||
|
||||
// 填充状态信息
|
||||
respList.forEach(resp -> {
|
||||
if (resp.getState() != null) {
|
||||
DeviceStateEnum stateEnum = DeviceStateEnum.valueOf(resp.getState());
|
||||
if (stateEnum != null) {
|
||||
resp.setStateName(stateEnum.getName());
|
||||
resp.setOnline(DeviceStateEnum.isOnline(resp.getState()));
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
return success(respList);
|
||||
}
|
||||
|
||||
@GetMapping("/statistics")
|
||||
@Operation(summary = "获取设备状态统计")
|
||||
@PreAuthorize("@ss.hasPermission('farm:device-status:query')")
|
||||
public CommonResult<DeviceStatusStatisticsRespVO> getDeviceStatusStatistics() {
|
||||
DeviceStatusStatisticsRespVO statistics = new DeviceStatusStatisticsRespVO();
|
||||
|
||||
// 统计各状态设备数量
|
||||
Long totalCount = deviceInfoMapper.selectCount();
|
||||
Long onlineCount = deviceInfoMapper.selectCountByState(DeviceStateEnum.ONLINE.getState());
|
||||
Long offlineCount = deviceInfoMapper.selectCountByState(DeviceStateEnum.OFFLINE.getState());
|
||||
Long inactiveCount = deviceInfoMapper.selectCountByState(DeviceStateEnum.INACTIVE.getState());
|
||||
Long maintenanceCount = deviceInfoMapper.selectCountByState(DeviceStateEnum.MAINTENANCE.getState());
|
||||
Long faultCount = deviceInfoMapper.selectCountByState(DeviceStateEnum.FAULT.getState());
|
||||
|
||||
statistics.setTotalCount(totalCount);
|
||||
statistics.setOnlineCount(onlineCount != null ? onlineCount : 0L);
|
||||
statistics.setOfflineCount(offlineCount != null ? offlineCount : 0L);
|
||||
statistics.setInactiveCount(inactiveCount != null ? inactiveCount : 0L);
|
||||
statistics.setMaintenanceCount(maintenanceCount != null ? maintenanceCount : 0L);
|
||||
statistics.setFaultCount(faultCount != null ? faultCount : 0L);
|
||||
|
||||
// 计算在线率和离线率
|
||||
if (totalCount > 0) {
|
||||
statistics.setOnlineRate((double) statistics.getOnlineCount() / totalCount * 100);
|
||||
statistics.setOfflineRate((double) statistics.getOfflineCount() / totalCount * 100);
|
||||
} else {
|
||||
statistics.setOnlineRate(0.0);
|
||||
statistics.setOfflineRate(0.0);
|
||||
}
|
||||
|
||||
return success(statistics);
|
||||
}
|
||||
|
||||
@PostMapping("/update-state/{deviceUniqueId}")
|
||||
@Operation(summary = "手动更新设备状态")
|
||||
@Parameter(name = "deviceUniqueId", description = "设备唯一标识", required = true, example = "2409AA01")
|
||||
@Parameter(name = "state", description = "设备状态", required = true, example = "1")
|
||||
@PreAuthorize("@ss.hasPermission('farm:device-status:update')")
|
||||
public CommonResult<Boolean> updateDeviceState(
|
||||
@PathVariable("deviceUniqueId") String deviceUniqueId,
|
||||
@RequestParam("state") Integer state) {
|
||||
deviceStatusService.updateDeviceState(deviceUniqueId, state);
|
||||
return success(true);
|
||||
}
|
||||
|
||||
@PostMapping("/set-online/{deviceUniqueId}")
|
||||
@Operation(summary = "手动设置设备上线")
|
||||
@Parameter(name = "deviceUniqueId", description = "设备唯一标识", required = true, example = "2409AA01")
|
||||
@PreAuthorize("@ss.hasPermission('farm:device-status:update')")
|
||||
public CommonResult<Boolean> setDeviceOnline(@PathVariable("deviceUniqueId") String deviceUniqueId) {
|
||||
deviceStatusService.deviceOnline(deviceUniqueId);
|
||||
return success(true);
|
||||
}
|
||||
|
||||
@PostMapping("/set-offline/{deviceUniqueId}")
|
||||
@Operation(summary = "手动设置设备离线")
|
||||
@Parameter(name = "deviceUniqueId", description = "设备唯一标识", required = true, example = "2409AA01")
|
||||
@PreAuthorize("@ss.hasPermission('farm:device-status:update')")
|
||||
public CommonResult<Boolean> setDeviceOffline(@PathVariable("deviceUniqueId") String deviceUniqueId) {
|
||||
deviceStatusService.deviceOffline(deviceUniqueId);
|
||||
return success(true);
|
||||
}
|
||||
|
||||
@PostMapping("/check-offline")
|
||||
@Operation(summary = "手动检查离线设备")
|
||||
@Parameter(name = "timeoutMinutes", description = "超时分钟数", required = false, example = "10")
|
||||
@PreAuthorize("@ss.hasPermission('farm:device-status:update')")
|
||||
public CommonResult<Integer> checkOfflineDevices(
|
||||
@RequestParam(value = "timeoutMinutes", defaultValue = "10") int timeoutMinutes) {
|
||||
int count = deviceStatusService.checkAndUpdateOfflineDevices(timeoutMinutes);
|
||||
return success(count);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,41 @@
|
||||
package cn.iocoder.yudao.module.farm.controller.admin.devicestatus.vo;
|
||||
|
||||
import cn.iocoder.yudao.framework.common.pojo.PageParam;
|
||||
import io.swagger.v3.oas.annotations.media.Schema;
|
||||
import lombok.Data;
|
||||
import lombok.EqualsAndHashCode;
|
||||
import lombok.ToString;
|
||||
import org.springframework.format.annotation.DateTimeFormat;
|
||||
|
||||
import java.time.LocalDateTime;
|
||||
|
||||
import static cn.iocoder.yudao.framework.common.util.date.DateUtils.FORMAT_YEAR_MONTH_DAY_HOUR_MINUTE_SECOND;
|
||||
|
||||
@Schema(description = "管理后台 - 设备状态分页 Request VO")
|
||||
@Data
|
||||
@EqualsAndHashCode(callSuper = true)
|
||||
@ToString(callSuper = true)
|
||||
public class DeviceStatusPageReqVO extends PageParam {
|
||||
|
||||
@Schema(description = "设备名称", example = "温湿度传感器")
|
||||
private String deviceName;
|
||||
|
||||
@Schema(description = "设备状态", example = "1")
|
||||
private Integer state;
|
||||
|
||||
@Schema(description = "设备分类名称", example = "传感器")
|
||||
private String deviceCategoryName;
|
||||
|
||||
@Schema(description = "基地名称", example = "示范基地1")
|
||||
private String baseName;
|
||||
|
||||
@Schema(description = "地块名称", example = "A1地块")
|
||||
private String plotName;
|
||||
|
||||
@Schema(description = "是否在线", example = "true")
|
||||
private Boolean online;
|
||||
|
||||
@Schema(description = "最后消息时间范围")
|
||||
@DateTimeFormat(pattern = FORMAT_YEAR_MONTH_DAY_HOUR_MINUTE_SECOND)
|
||||
private LocalDateTime[] lastMessageTime;
|
||||
}
|
||||
@@ -0,0 +1,47 @@
|
||||
package cn.iocoder.yudao.module.farm.controller.admin.devicestatus.vo;
|
||||
|
||||
import io.swagger.v3.oas.annotations.media.Schema;
|
||||
import lombok.Data;
|
||||
|
||||
import java.time.LocalDateTime;
|
||||
|
||||
@Schema(description = "管理后台 - 设备状态 Response VO")
|
||||
@Data
|
||||
public class DeviceStatusRespVO {
|
||||
|
||||
@Schema(description = "设备唯一标识", requiredMode = Schema.RequiredMode.REQUIRED, example = "2409AA01")
|
||||
private String deviceUniqueId;
|
||||
|
||||
@Schema(description = "设备名称", example = "温湿度传感器1")
|
||||
private String deviceName;
|
||||
|
||||
@Schema(description = "设备状态", requiredMode = Schema.RequiredMode.REQUIRED, example = "1")
|
||||
private Integer state;
|
||||
|
||||
@Schema(description = "设备状态名称", example = "在线")
|
||||
private String stateName;
|
||||
|
||||
@Schema(description = "是否在线", example = "true")
|
||||
private Boolean online;
|
||||
|
||||
@Schema(description = "设备激活时间")
|
||||
private LocalDateTime activeTime;
|
||||
|
||||
@Schema(description = "设备最后上线时间")
|
||||
private LocalDateTime onlineTime;
|
||||
|
||||
@Schema(description = "设备最后离线时间")
|
||||
private LocalDateTime offlineTime;
|
||||
|
||||
@Schema(description = "设备最后消息时间")
|
||||
private LocalDateTime lastMessageTime;
|
||||
|
||||
@Schema(description = "设备分类名称", example = "传感器")
|
||||
private String deviceCategoryName;
|
||||
|
||||
@Schema(description = "基地名称", example = "示范基地1")
|
||||
private String baseName;
|
||||
|
||||
@Schema(description = "地块名称", example = "A1地块")
|
||||
private String plotName;
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
package cn.iocoder.yudao.module.farm.controller.admin.devicestatus.vo;
|
||||
|
||||
import io.swagger.v3.oas.annotations.media.Schema;
|
||||
import lombok.Data;
|
||||
|
||||
@Schema(description = "管理后台 - 设备状态统计 Response VO")
|
||||
@Data
|
||||
public class DeviceStatusStatisticsRespVO {
|
||||
|
||||
@Schema(description = "设备总数", example = "100")
|
||||
private Long totalCount;
|
||||
|
||||
@Schema(description = "在线设备数", example = "85")
|
||||
private Long onlineCount;
|
||||
|
||||
@Schema(description = "离线设备数", example = "10")
|
||||
private Long offlineCount;
|
||||
|
||||
@Schema(description = "未激活设备数", example = "3")
|
||||
private Long inactiveCount;
|
||||
|
||||
@Schema(description = "维护中设备数", example = "1")
|
||||
private Long maintenanceCount;
|
||||
|
||||
@Schema(description = "故障设备数", example = "1")
|
||||
private Long faultCount;
|
||||
|
||||
@Schema(description = "在线率(百分比)", example = "85.0")
|
||||
private Double onlineRate;
|
||||
|
||||
@Schema(description = "离线率(百分比)", example = "10.0")
|
||||
private Double offlineRate;
|
||||
}
|
||||
@@ -0,0 +1,121 @@
|
||||
package cn.iocoder.yudao.module.farm.controller.admin.mqtt;
|
||||
|
||||
import org.springframework.web.bind.annotation.*;
|
||||
import jakarta.annotation.Resource;
|
||||
import org.springframework.validation.annotation.Validated;
|
||||
import org.springframework.security.access.prepost.PreAuthorize;
|
||||
import io.swagger.v3.oas.annotations.tags.Tag;
|
||||
import io.swagger.v3.oas.annotations.Parameter;
|
||||
import io.swagger.v3.oas.annotations.Operation;
|
||||
|
||||
import jakarta.validation.constraints.*;
|
||||
import jakarta.validation.Valid;
|
||||
import jakarta.servlet.http.*;
|
||||
|
||||
import java.util.*;
|
||||
import java.io.IOException;
|
||||
|
||||
import cn.iocoder.yudao.framework.common.pojo.PageParam;
|
||||
import cn.iocoder.yudao.framework.common.pojo.PageResult;
|
||||
import cn.iocoder.yudao.framework.common.pojo.CommonResult;
|
||||
import cn.iocoder.yudao.framework.common.util.object.BeanUtils;
|
||||
|
||||
import static cn.iocoder.yudao.framework.common.pojo.CommonResult.success;
|
||||
|
||||
import cn.iocoder.yudao.framework.excel.core.util.ExcelUtils;
|
||||
|
||||
import cn.iocoder.yudao.framework.apilog.core.annotation.ApiAccessLog;
|
||||
|
||||
import static cn.iocoder.yudao.framework.apilog.core.enums.OperateTypeEnum.*;
|
||||
|
||||
import cn.iocoder.yudao.module.farm.controller.admin.mqtt.vo.*;
|
||||
import cn.iocoder.yudao.module.farm.dal.dataobject.mqtt.MqttConnectionDO;
|
||||
import cn.iocoder.yudao.module.farm.service.mqtt.MqttConnectionManagementService;
|
||||
|
||||
@Tag(name = "管理后台 - MQTT连接配置")
|
||||
@RestController
|
||||
@RequestMapping("/farm/mqtt-connection")
|
||||
@Validated
|
||||
public class MqttConnectionController {
|
||||
|
||||
@Resource
|
||||
private MqttConnectionManagementService mqttConnectionManagementService;
|
||||
|
||||
@PostMapping("/create")
|
||||
@Operation(summary = "创建MQTT连接配置")
|
||||
@PreAuthorize("@ss.hasPermission('farm:mqtt-connection:create')")
|
||||
public CommonResult<Long> createMqttConnection(@Valid @RequestBody MqttConnectionSaveReqVO createReqVO) {
|
||||
return success(mqttConnectionManagementService.createMqttConnection(createReqVO));
|
||||
}
|
||||
|
||||
@PutMapping("/update")
|
||||
@Operation(summary = "更新MQTT连接配置")
|
||||
@PreAuthorize("@ss.hasPermission('farm:mqtt-connection:update')")
|
||||
public CommonResult<Boolean> updateMqttConnection(@Valid @RequestBody MqttConnectionSaveReqVO updateReqVO) {
|
||||
mqttConnectionManagementService.updateMqttConnection(updateReqVO);
|
||||
return success(true);
|
||||
}
|
||||
|
||||
@DeleteMapping("/delete")
|
||||
@Operation(summary = "删除MQTT连接配置")
|
||||
@Parameter(name = "id", description = "编号", required = true)
|
||||
@PreAuthorize("@ss.hasPermission('farm:mqtt-connection:delete')")
|
||||
public CommonResult<Boolean> deleteMqttConnection(@RequestParam("id") Long id) {
|
||||
mqttConnectionManagementService.deleteMqttConnection(id);
|
||||
return success(true);
|
||||
}
|
||||
|
||||
@DeleteMapping("/delete-batch")
|
||||
@Operation(summary = "批量删除MQTT连接配置")
|
||||
@Parameter(name = "ids", description = "编号列表", required = true)
|
||||
@PreAuthorize("@ss.hasPermission('farm:mqtt-connection:delete')")
|
||||
public CommonResult<Boolean> deleteMqttConnectionListByIds(@RequestParam("ids") List<Long> ids) {
|
||||
mqttConnectionManagementService.deleteMqttConnectionListByIds(ids);
|
||||
return success(true);
|
||||
}
|
||||
|
||||
@GetMapping("/get")
|
||||
@Operation(summary = "获得MQTT连接配置")
|
||||
@Parameter(name = "id", description = "编号", required = true, example = "1024")
|
||||
@PreAuthorize("@ss.hasPermission('farm:mqtt-connection:query')")
|
||||
public CommonResult<MqttConnectionRespVO> getMqttConnection(@RequestParam("id") Long id) {
|
||||
MqttConnectionDO mqttConnection = mqttConnectionManagementService.getMqttConnection(id);
|
||||
return success(BeanUtils.toBean(mqttConnection, MqttConnectionRespVO.class));
|
||||
}
|
||||
|
||||
@GetMapping("/page")
|
||||
@Operation(summary = "获得MQTT连接配置分页")
|
||||
@PreAuthorize("@ss.hasPermission('farm:mqtt-connection:query')")
|
||||
public CommonResult<PageResult<MqttConnectionRespVO>> getMqttConnectionPage(@Valid MqttConnectionPageReqVO pageReqVO) {
|
||||
PageResult<MqttConnectionDO> pageResult = mqttConnectionManagementService.getMqttConnectionPage(pageReqVO);
|
||||
return success(BeanUtils.toBean(pageResult, MqttConnectionRespVO.class));
|
||||
}
|
||||
|
||||
|
||||
@GetMapping("/get-by-device/{deviceId}")
|
||||
@Operation(summary = "根据设备ID获取MQTT连接配置")
|
||||
@Parameter(name = "deviceId", description = "设备ID", required = true, example = "2409AA01")
|
||||
@PreAuthorize("@ss.hasPermission('farm:mqtt-connection:query')")
|
||||
public CommonResult<MqttConnectionRespVO> getMqttConnectionByDeviceId(@PathVariable("deviceId") String deviceId) {
|
||||
MqttConnectionDO mqttConnection = mqttConnectionManagementService.getMqttConnectionByDeviceId(deviceId);
|
||||
return success(BeanUtils.toBean(mqttConnection, MqttConnectionRespVO.class));
|
||||
}
|
||||
|
||||
@PostMapping("/test/{id}")
|
||||
@Operation(summary = "测试MQTT连接")
|
||||
@Parameter(name = "id", description = "连接配置ID", required = true, example = "1024")
|
||||
@PreAuthorize("@ss.hasPermission('farm:mqtt-connection:test')")
|
||||
public CommonResult<Boolean> testMqttConnection(@PathVariable("id") Long id) {
|
||||
boolean result = mqttConnectionManagementService.testMqttConnection(id);
|
||||
return success(result);
|
||||
}
|
||||
|
||||
@PostMapping("/reload/{deviceId}")
|
||||
@Operation(summary = "重新加载MQTT连接")
|
||||
@Parameter(name = "deviceId", description = "设备ID", required = true, example = "2409AA01")
|
||||
@PreAuthorize("@ss.hasPermission('farm:mqtt-connection:reload')")
|
||||
public CommonResult<Boolean> reloadMqttConnection(@PathVariable("deviceId") String deviceId) {
|
||||
mqttConnectionManagementService.reloadMqttConnection(deviceId);
|
||||
return success(true);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,128 @@
|
||||
package cn.iocoder.yudao.module.farm.controller.admin.mqtt;
|
||||
|
||||
import cn.iocoder.yudao.framework.common.pojo.CommonResult;
|
||||
import cn.iocoder.yudao.module.farm.service.mqtt.DeviceMqttService;
|
||||
import cn.iocoder.yudao.module.farm.service.mqtt.MqttConnectionService;
|
||||
import cn.iocoder.yudao.module.farm.service.mqtt.MqttDeviceInitService;
|
||||
import io.swagger.v3.oas.annotations.Operation;
|
||||
import io.swagger.v3.oas.annotations.Parameter;
|
||||
import io.swagger.v3.oas.annotations.tags.Tag;
|
||||
import jakarta.annotation.Resource;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.eclipse.paho.client.mqttv3.MqttClient;
|
||||
import org.eclipse.paho.client.mqttv3.MqttException;
|
||||
import org.springframework.security.access.prepost.PreAuthorize;
|
||||
import org.springframework.validation.annotation.Validated;
|
||||
import org.springframework.web.bind.annotation.*;
|
||||
|
||||
import static cn.iocoder.yudao.framework.common.pojo.CommonResult.success;
|
||||
|
||||
@Tag(name = "管理后台 - MQTT测试")
|
||||
@RestController
|
||||
@RequestMapping("/farm/mqtt-test")
|
||||
@Validated
|
||||
@Slf4j
|
||||
public class MqttTestController {
|
||||
|
||||
@Resource
|
||||
private DeviceMqttService deviceMqttService;
|
||||
|
||||
@Resource
|
||||
private MqttConnectionService mqttConnectionService;
|
||||
|
||||
@Resource
|
||||
private MqttDeviceInitService mqttDeviceInitService;
|
||||
|
||||
@PostMapping("/init-device/{deviceId}")
|
||||
@Operation(summary = "初始化设备MQTT连接配置")
|
||||
@Parameter(name = "deviceId", description = "设备ID", required = true, example = "2409AA01")
|
||||
@PreAuthorize("@ss.hasPermission('farm:mqtt-connection:create')")
|
||||
public CommonResult<Boolean> initDeviceConnection(
|
||||
@PathVariable("deviceId") String deviceId,
|
||||
@RequestParam(value = "tenantId", defaultValue = "1") String tenantId,
|
||||
@RequestParam(value = "broker", defaultValue = "tcp://mqtt.api.mimirii.com:1883") String broker,
|
||||
@RequestParam(value = "username", defaultValue = "mimir_user") String username,
|
||||
@RequestParam(value = "password", defaultValue = "mimir123456@") String password,
|
||||
@RequestParam(value = "telemetryTopic", defaultValue = "/WFM/{deviceId}/up") String telemetryTopic,
|
||||
@RequestParam(value = "commandTopic", defaultValue = "/WFM/{deviceId}/down") String commandTopic) {
|
||||
|
||||
mqttDeviceInitService.initDeviceConnection(deviceId, tenantId, broker, username, password, telemetryTopic, commandTopic);
|
||||
return success(true);
|
||||
}
|
||||
|
||||
@PostMapping("/test-connection/{deviceId}")
|
||||
@Operation(summary = "测试设备MQTT连接")
|
||||
@Parameter(name = "deviceId", description = "设备ID", required = true, example = "2409AA01")
|
||||
@PreAuthorize("@ss.hasPermission('farm:mqtt-connection:test')")
|
||||
public CommonResult<String> testConnection(@PathVariable("deviceId") String deviceId) {
|
||||
try {
|
||||
MqttClient client = mqttConnectionService.getClientByDeviceId(deviceId);
|
||||
if (client == null) {
|
||||
return success("连接失败:未找到设备连接配置");
|
||||
}
|
||||
if (client.isConnected()) {
|
||||
return success("连接成功:设备已连接到MQTT服务器");
|
||||
} else {
|
||||
return success("连接失败:设备未连接到MQTT服务器");
|
||||
}
|
||||
} catch (MqttException e) {
|
||||
log.error("测试MQTT连接失败", e);
|
||||
return success("连接失败:" + e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
@PostMapping("/publish-test/{deviceId}")
|
||||
@Operation(summary = "发送测试消息")
|
||||
@Parameter(name = "deviceId", description = "设备ID", required = true, example = "2409AA01")
|
||||
@PreAuthorize("@ss.hasPermission('farm:mqtt-connection:test')")
|
||||
public CommonResult<String> publishTest(
|
||||
@PathVariable("deviceId") String deviceId,
|
||||
@RequestParam(value = "message", defaultValue = "{\"temperature\":25.5,\"humidity\":60.2}") String message) {
|
||||
try {
|
||||
deviceMqttService.publishTelemetry(deviceId, message);
|
||||
return success("消息发送成功");
|
||||
} catch (MqttException e) {
|
||||
log.error("发送测试消息失败", e);
|
||||
return success("消息发送失败:" + e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
@PostMapping("/subscribe-test/{deviceId}")
|
||||
@Operation(summary = "订阅测试")
|
||||
@Parameter(name = "deviceId", description = "设备ID", required = true, example = "2409AA01")
|
||||
@PreAuthorize("@ss.hasPermission('farm:mqtt-connection:test')")
|
||||
public CommonResult<String> subscribeTest(@PathVariable("deviceId") String deviceId) {
|
||||
try {
|
||||
deviceMqttService.subscribeTelemetryAndPersist(deviceId);
|
||||
return success("订阅成功:已开始监听设备上报数据");
|
||||
} catch (MqttException e) {
|
||||
log.error("订阅测试失败", e);
|
||||
return success("订阅失败:" + e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
@PostMapping("/send-command/{deviceId}")
|
||||
@Operation(summary = "发送设备指令")
|
||||
@Parameter(name = "deviceId", description = "设备ID", required = true, example = "2409AA01")
|
||||
@PreAuthorize("@ss.hasPermission('farm:mqtt-connection:test')")
|
||||
public CommonResult<String> sendCommand(
|
||||
@PathVariable("deviceId") String deviceId,
|
||||
@RequestParam(value = "command", defaultValue = "{\"PUMP\":\"1\"}") String command) {
|
||||
try {
|
||||
deviceMqttService.sendCommand(deviceId, command);
|
||||
return success("指令发送成功");
|
||||
} catch (MqttException e) {
|
||||
log.error("发送设备指令失败", e);
|
||||
return success("指令发送失败:" + e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
@PostMapping("/reload-connection/{deviceId}")
|
||||
@Operation(summary = "重新加载设备连接")
|
||||
@Parameter(name = "deviceId", description = "设备ID", required = true, example = "2409AA01")
|
||||
@PreAuthorize("@ss.hasPermission('farm:mqtt-connection:reload')")
|
||||
public CommonResult<String> reloadConnection(@PathVariable("deviceId") String deviceId) {
|
||||
mqttConnectionService.reload(deviceId);
|
||||
return success("连接已重新加载");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,32 @@
|
||||
package cn.iocoder.yudao.module.farm.controller.admin.mqtt.vo;
|
||||
|
||||
import lombok.*;
|
||||
import io.swagger.v3.oas.annotations.media.Schema;
|
||||
import cn.iocoder.yudao.framework.common.pojo.PageParam;
|
||||
import org.springframework.format.annotation.DateTimeFormat;
|
||||
import java.time.LocalDateTime;
|
||||
|
||||
import static cn.iocoder.yudao.framework.common.util.date.DateUtils.FORMAT_YEAR_MONTH_DAY_HOUR_MINUTE_SECOND;
|
||||
|
||||
@Schema(description = "管理后台 - MQTT连接配置分页 Request VO")
|
||||
@Data
|
||||
@EqualsAndHashCode(callSuper = true)
|
||||
@ToString(callSuper = true)
|
||||
public class MqttConnectionPageReqVO extends PageParam {
|
||||
|
||||
@Schema(description = "连接名称", example = "设备连接1")
|
||||
private String name;
|
||||
|
||||
@Schema(description = "设备ID", example = "2409AA01")
|
||||
private String deviceId;
|
||||
|
||||
@Schema(description = "MQTT服务器地址", example = "tcp://mqtt.api.mimirii.com:1883")
|
||||
private String broker;
|
||||
|
||||
@Schema(description = "连接状态", example = "true")
|
||||
private Boolean status;
|
||||
|
||||
@Schema(description = "创建时间")
|
||||
@DateTimeFormat(pattern = FORMAT_YEAR_MONTH_DAY_HOUR_MINUTE_SECOND)
|
||||
private LocalDateTime[] createTime;
|
||||
}
|
||||
@@ -0,0 +1,62 @@
|
||||
package cn.iocoder.yudao.module.farm.controller.admin.mqtt.vo;
|
||||
|
||||
import cn.idev.excel.annotation.ExcelIgnoreUnannotated;
|
||||
import cn.idev.excel.annotation.ExcelProperty;
|
||||
import io.swagger.v3.oas.annotations.media.Schema;
|
||||
import lombok.Data;
|
||||
|
||||
import java.time.LocalDateTime;
|
||||
|
||||
@Schema(description = "管理后台 - MQTT连接配置 Response VO")
|
||||
@Data
|
||||
@ExcelIgnoreUnannotated
|
||||
public class MqttConnectionRespVO {
|
||||
|
||||
@Schema(description = "编号", requiredMode = Schema.RequiredMode.REQUIRED, example = "1024")
|
||||
private Long id;
|
||||
|
||||
@Schema(description = "租户ID", requiredMode = Schema.RequiredMode.REQUIRED, example = "1")
|
||||
private String tenantId;
|
||||
|
||||
@Schema(description = "连接名称", requiredMode = Schema.RequiredMode.REQUIRED, example = "设备连接1")
|
||||
private String name;
|
||||
|
||||
@Schema(description = "设备ID", requiredMode = Schema.RequiredMode.REQUIRED, example = "2409AA01")
|
||||
private String deviceId;
|
||||
|
||||
@Schema(description = "MQTT服务器地址", requiredMode = Schema.RequiredMode.REQUIRED, example = "tcp://mqtt.api.mimirii.com:1883")
|
||||
private String broker;
|
||||
|
||||
@Schema(description = "客户端ID", requiredMode = Schema.RequiredMode.REQUIRED, example = "2409AA01")
|
||||
private String clientId;
|
||||
|
||||
@Schema(description = "用户名", example = "mimir_user")
|
||||
private String username;
|
||||
|
||||
@Schema(description = "密码", example = "mimir123456@")
|
||||
private String password;
|
||||
|
||||
@Schema(description = "清除会话", example = "true")
|
||||
private Boolean cleanSession;
|
||||
|
||||
@Schema(description = "保持连接间隔(秒)", example = "30")
|
||||
private Integer keepAlive;
|
||||
|
||||
@Schema(description = "消息质量等级", example = "1")
|
||||
private Integer qos;
|
||||
|
||||
@Schema(description = "自动重连", example = "true")
|
||||
private Boolean reconnect;
|
||||
|
||||
@Schema(description = "连接状态", example = "true")
|
||||
private Boolean status;
|
||||
|
||||
@Schema(description = "设备上报主题模板", example = "/WFM/{deviceId}/up")
|
||||
private String telemetryTopic;
|
||||
|
||||
@Schema(description = "设备下发主题模板", example = "/WFM/{deviceId}/down")
|
||||
private String commandTopic;
|
||||
|
||||
@Schema(description = "创建时间", requiredMode = Schema.RequiredMode.REQUIRED)
|
||||
private LocalDateTime createTime;
|
||||
}
|
||||
@@ -0,0 +1,58 @@
|
||||
package cn.iocoder.yudao.module.farm.controller.admin.mqtt.vo;
|
||||
|
||||
import io.swagger.v3.oas.annotations.media.Schema;
|
||||
import lombok.Data;
|
||||
|
||||
import jakarta.validation.constraints.NotEmpty;
|
||||
import jakarta.validation.constraints.NotNull;
|
||||
|
||||
@Schema(description = "管理后台 - MQTT连接配置新增/修改 Request VO")
|
||||
@Data
|
||||
public class MqttConnectionSaveReqVO {
|
||||
|
||||
@Schema(description = "编号", requiredMode = Schema.RequiredMode.REQUIRED, example = "1024")
|
||||
private Long id;
|
||||
|
||||
@Schema(description = "连接名称", requiredMode = Schema.RequiredMode.REQUIRED, example = "设备连接1")
|
||||
@NotEmpty(message = "连接名称不能为空")
|
||||
private String name;
|
||||
|
||||
@Schema(description = "设备ID", requiredMode = Schema.RequiredMode.REQUIRED, example = "2409AA01")
|
||||
@NotEmpty(message = "设备ID不能为空")
|
||||
private String deviceId;
|
||||
|
||||
@Schema(description = "MQTT服务器地址", requiredMode = Schema.RequiredMode.REQUIRED, example = "tcp://mqtt.api.mimirii.com:1883")
|
||||
@NotEmpty(message = "MQTT服务器地址不能为空")
|
||||
private String broker;
|
||||
|
||||
@Schema(description = "客户端ID", requiredMode = Schema.RequiredMode.REQUIRED, example = "2409AA01")
|
||||
@NotEmpty(message = "客户端ID不能为空")
|
||||
private String clientId;
|
||||
|
||||
@Schema(description = "用户名", example = "mimir_user")
|
||||
private String username;
|
||||
|
||||
@Schema(description = "密码", example = "mimir123456@")
|
||||
private String password;
|
||||
|
||||
@Schema(description = "清除会话", example = "true")
|
||||
private Boolean cleanSession = true;
|
||||
|
||||
@Schema(description = "保持连接间隔(秒)", example = "30")
|
||||
private Integer keepAlive = 30;
|
||||
|
||||
@Schema(description = "消息质量等级", example = "1")
|
||||
private Integer qos = 1;
|
||||
|
||||
@Schema(description = "自动重连", example = "true")
|
||||
private Boolean reconnect = true;
|
||||
|
||||
@Schema(description = "连接状态", example = "true")
|
||||
private Boolean status = true;
|
||||
|
||||
@Schema(description = "设备上报主题模板", example = "/WFM/{deviceId}/up")
|
||||
private String telemetryTopic;
|
||||
|
||||
@Schema(description = "设备下发主题模板", example = "/WFM/{deviceId}/down")
|
||||
private String commandTopic;
|
||||
}
|
||||
@@ -2,9 +2,6 @@ package cn.iocoder.yudao.module.farm.dal.dataobject.deviceinfo;
|
||||
|
||||
import lombok.*;
|
||||
|
||||
import java.time.LocalDateTime;
|
||||
import java.time.LocalDateTime;
|
||||
import java.time.LocalDateTime;
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.List;
|
||||
|
||||
@@ -166,6 +163,31 @@ public class DeviceInfoDO extends BaseDO {
|
||||
@TableField(value = "data_image")
|
||||
private String dataImage;
|
||||
|
||||
/**
|
||||
* 设备状态:0-未激活,1-在线,2-离线,3-维护中,4-故障
|
||||
*/
|
||||
private Integer state;
|
||||
|
||||
/**
|
||||
* 设备激活时间
|
||||
*/
|
||||
private LocalDateTime activeTime;
|
||||
|
||||
/**
|
||||
* 设备最后上线时间
|
||||
*/
|
||||
private LocalDateTime onlineTime;
|
||||
|
||||
/**
|
||||
* 设备最后离线时间
|
||||
*/
|
||||
private LocalDateTime offlineTime;
|
||||
|
||||
/**
|
||||
* 设备最后消息时间
|
||||
*/
|
||||
private LocalDateTime lastMessageTime;
|
||||
|
||||
/**
|
||||
* 数据快照集合(业务层使用)
|
||||
*/
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
package cn.iocoder.yudao.module.farm.dal.dataobject.modelenvironmentindex;
|
||||
|
||||
import cn.iocoder.yudao.framework.excel.core.annotations.DictFormat;
|
||||
import lombok.*;
|
||||
|
||||
import java.math.BigDecimal;
|
||||
@@ -36,6 +37,7 @@ public class ModelEnvironmentIndexDO extends BaseDO {
|
||||
/**
|
||||
* 设备分类
|
||||
*/
|
||||
@DictFormat("device_class")
|
||||
private String deviceClass;
|
||||
/**
|
||||
* 环境指标
|
||||
|
||||
@@ -42,6 +42,16 @@ public class MqttConnectionDO extends BaseDO {
|
||||
private Boolean reconnect;
|
||||
|
||||
private Boolean status;
|
||||
|
||||
/**
|
||||
* 设备上报主题模板 (如: /WFM/{deviceId}/up)
|
||||
*/
|
||||
private String telemetryTopic;
|
||||
|
||||
/**
|
||||
* 设备下发主题模板 (如: /WFM/{deviceId}/down)
|
||||
*/
|
||||
private String commandTopic;
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -6,6 +6,10 @@ import cn.iocoder.yudao.framework.mybatis.core.mapper.BaseMapperX;
|
||||
import cn.iocoder.yudao.module.farm.controller.admin.deviceinfo.vo.DeviceInfoPageReqVO;
|
||||
import cn.iocoder.yudao.module.farm.dal.dataobject.deviceinfo.DeviceInfoDO;
|
||||
import org.apache.ibatis.annotations.Mapper;
|
||||
import org.apache.ibatis.annotations.Param;
|
||||
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* 设备信息 Mapper
|
||||
@@ -53,4 +57,33 @@ public interface DeviceInfoMapper extends BaseMapperX<DeviceInfoDO> {
|
||||
.orderByDesc(DeviceInfoDO::getCreateTime));
|
||||
}
|
||||
|
||||
/**
|
||||
* 根据设备状态查询设备列表
|
||||
*/
|
||||
default List<DeviceInfoDO> selectListByState(Integer state) {
|
||||
return selectList(new LambdaQueryWrapperX<DeviceInfoDO>()
|
||||
.eq(DeviceInfoDO::getState, state));
|
||||
}
|
||||
|
||||
/**
|
||||
* 查询超时未上报的设备列表
|
||||
*/
|
||||
default List<DeviceInfoDO> selectTimeoutDevices(@Param("timeoutTime") LocalDateTime timeoutTime) {
|
||||
return selectList(new LambdaQueryWrapperX<DeviceInfoDO>()
|
||||
.eq(DeviceInfoDO::getState, 1) // 只查询在线状态的设备
|
||||
.and(wrapper -> wrapper
|
||||
.isNull(DeviceInfoDO::getLastMessageTime) // 从未上报过消息
|
||||
.or()
|
||||
.lt(DeviceInfoDO::getLastMessageTime, timeoutTime) // 或者最后消息时间超时
|
||||
));
|
||||
}
|
||||
|
||||
/**
|
||||
* 统计各状态设备数量
|
||||
*/
|
||||
default Long selectCountByState(Integer state) {
|
||||
return selectCount(new LambdaQueryWrapperX<DeviceInfoDO>()
|
||||
.eq(DeviceInfoDO::getState, state));
|
||||
}
|
||||
|
||||
}
|
||||
@@ -1,7 +1,9 @@
|
||||
package cn.iocoder.yudao.module.farm.dal.mysql.mqtt;
|
||||
|
||||
import cn.iocoder.yudao.framework.common.pojo.PageResult;
|
||||
import cn.iocoder.yudao.framework.mybatis.core.mapper.BaseMapperX;
|
||||
import cn.iocoder.yudao.framework.mybatis.core.query.LambdaQueryWrapperX;
|
||||
import cn.iocoder.yudao.module.farm.controller.admin.mqtt.vo.MqttConnectionPageReqVO;
|
||||
import cn.iocoder.yudao.module.farm.dal.dataobject.mqtt.MqttConnectionDO;
|
||||
import org.apache.ibatis.annotations.Mapper;
|
||||
|
||||
@@ -15,6 +17,16 @@ public interface MqttConnectionMapper extends BaseMapperX<MqttConnectionDO> {
|
||||
.eq(MqttConnectionDO::getStatus, true)
|
||||
);
|
||||
}
|
||||
|
||||
default PageResult<MqttConnectionDO> selectPage(MqttConnectionPageReqVO reqVO) {
|
||||
return selectPage(reqVO, new LambdaQueryWrapperX<MqttConnectionDO>()
|
||||
.likeIfPresent(MqttConnectionDO::getName, reqVO.getName())
|
||||
.eqIfPresent(MqttConnectionDO::getDeviceId, reqVO.getDeviceId())
|
||||
.likeIfPresent(MqttConnectionDO::getBroker, reqVO.getBroker())
|
||||
.eqIfPresent(MqttConnectionDO::getStatus, reqVO.getStatus())
|
||||
.betweenIfPresent(MqttConnectionDO::getCreateTime, reqVO.getCreateTime())
|
||||
.orderByDesc(MqttConnectionDO::getId));
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,63 @@
|
||||
package cn.iocoder.yudao.module.farm.enums;
|
||||
|
||||
import cn.iocoder.yudao.framework.common.core.ArrayValuable;
|
||||
import lombok.Getter;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
|
||||
import java.util.Arrays;
|
||||
|
||||
/**
|
||||
* 农场设备状态枚举
|
||||
*
|
||||
* @author 芋道源码
|
||||
*/
|
||||
@RequiredArgsConstructor
|
||||
@Getter
|
||||
public enum DeviceStateEnum implements ArrayValuable<Integer> {
|
||||
|
||||
INACTIVE(0, "未激活"),
|
||||
ONLINE(1, "在线"),
|
||||
OFFLINE(2, "离线"),
|
||||
MAINTENANCE(3, "维护中"),
|
||||
FAULT(4, "故障");
|
||||
|
||||
public static final Integer[] ARRAYS = Arrays.stream(values()).map(DeviceStateEnum::getState).toArray(Integer[]::new);
|
||||
|
||||
/**
|
||||
* 状态
|
||||
*/
|
||||
private final Integer state;
|
||||
/**
|
||||
* 状态名
|
||||
*/
|
||||
private final String name;
|
||||
|
||||
@Override
|
||||
public Integer[] array() {
|
||||
return ARRAYS;
|
||||
}
|
||||
|
||||
/**
|
||||
* 判断设备是否在线
|
||||
*/
|
||||
public static boolean isOnline(Integer state) {
|
||||
return ONLINE.getState().equals(state);
|
||||
}
|
||||
|
||||
/**
|
||||
* 判断设备是否离线
|
||||
*/
|
||||
public static boolean isOffline(Integer state) {
|
||||
return OFFLINE.getState().equals(state);
|
||||
}
|
||||
|
||||
/**
|
||||
* 根据状态值获取枚举
|
||||
*/
|
||||
public static DeviceStateEnum valueOf(Integer state) {
|
||||
return Arrays.stream(values())
|
||||
.filter(item -> item.getState().equals(state))
|
||||
.findFirst()
|
||||
.orElse(null);
|
||||
}
|
||||
}
|
||||
@@ -120,4 +120,7 @@ public interface ErrorCodeConstants {
|
||||
ErrorCode PLOT_AREA_ERRORS = new ErrorCode(1_050_036_000, "地块面积总和大于基地面积");
|
||||
|
||||
ErrorCode INDICATOR_NAME_NOT_EXISTS = new ErrorCode(1_050_037_000, "指标名称不存在");
|
||||
|
||||
// ========== MQTT连接配置 1_050_038_000 ==========
|
||||
ErrorCode MQTT_CONNECTION_NOT_EXISTS = new ErrorCode(1_050_038_000, "MQTT连接配置不存在");
|
||||
}
|
||||
@@ -0,0 +1,65 @@
|
||||
package cn.iocoder.yudao.module.farm.job;
|
||||
|
||||
import cn.hutool.json.JSONUtil;
|
||||
import cn.iocoder.yudao.framework.quartz.core.handler.JobHandler;
|
||||
import cn.iocoder.yudao.module.farm.service.devicestatus.DeviceStatusService;
|
||||
import jakarta.annotation.Resource;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
* 农场设备离线检查 Job
|
||||
* <p>
|
||||
* 检测逻辑:设备最后消息时间超过一定时间,则认为设备离线
|
||||
*
|
||||
* @author 芋道源码
|
||||
*/
|
||||
@Component
|
||||
@Slf4j
|
||||
public class DeviceOfflineCheckJob implements JobHandler {
|
||||
|
||||
/**
|
||||
* 设备离线超时时间(分钟)
|
||||
* <p>
|
||||
* 可以通过定时任务参数配置,默认10分钟
|
||||
*/
|
||||
public static final int DEFAULT_OFFLINE_TIMEOUT_MINUTES = 10;
|
||||
|
||||
@Resource
|
||||
private DeviceStatusService deviceStatusService;
|
||||
|
||||
@Override
|
||||
public String execute(String param) {
|
||||
log.info("开始执行设备离线检查任务,参数:{}", param);
|
||||
|
||||
// 解析参数,获取超时时间
|
||||
int timeoutMinutes = DEFAULT_OFFLINE_TIMEOUT_MINUTES;
|
||||
if (param != null && !param.trim().isEmpty()) {
|
||||
try {
|
||||
Map<String, Object> paramMap = JSONUtil.toBean(param, Map.class);
|
||||
if (paramMap.containsKey("timeoutMinutes")) {
|
||||
timeoutMinutes = Integer.parseInt(paramMap.get("timeoutMinutes").toString());
|
||||
}
|
||||
} catch (Exception e) {
|
||||
log.warn("解析任务参数失败,使用默认超时时间:{} 分钟", DEFAULT_OFFLINE_TIMEOUT_MINUTES);
|
||||
}
|
||||
}
|
||||
|
||||
// 检查并更新离线设备
|
||||
int offlineCount = deviceStatusService.checkAndUpdateOfflineDevices(timeoutMinutes);
|
||||
|
||||
// 构建返回结果
|
||||
Map<String, Object> result = new HashMap<>();
|
||||
result.put("timeoutMinutes", timeoutMinutes);
|
||||
result.put("offlineDeviceCount", offlineCount);
|
||||
result.put("executeTime", System.currentTimeMillis());
|
||||
|
||||
String resultJson = JSONUtil.toJsonStr(result);
|
||||
log.info("设备离线检查任务执行完成:{}", resultJson);
|
||||
|
||||
return resultJson;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,108 @@
|
||||
package cn.iocoder.yudao.module.farm.service.devicestatus;
|
||||
|
||||
import cn.iocoder.yudao.module.farm.dal.dataobject.deviceinfo.DeviceInfoDO;
|
||||
import cn.iocoder.yudao.module.farm.enums.DeviceStateEnum;
|
||||
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* 设备状态管理 Service 接口
|
||||
*
|
||||
* @author 芋道源码
|
||||
*/
|
||||
public interface DeviceStatusService {
|
||||
|
||||
/**
|
||||
* 更新设备状态
|
||||
*
|
||||
* @param deviceUniqueId 设备唯一标识
|
||||
* @param state 设备状态
|
||||
*/
|
||||
void updateDeviceState(String deviceUniqueId, Integer state);
|
||||
|
||||
/**
|
||||
* 更新设备状态
|
||||
*
|
||||
* @param deviceUniqueId 设备唯一标识
|
||||
* @param state 设备状态
|
||||
* @param messageTime 消息时间
|
||||
*/
|
||||
void updateDeviceState(String deviceUniqueId, Integer state, LocalDateTime messageTime);
|
||||
|
||||
/**
|
||||
* 设备上线
|
||||
*
|
||||
* @param deviceUniqueId 设备唯一标识
|
||||
*/
|
||||
void deviceOnline(String deviceUniqueId);
|
||||
|
||||
/**
|
||||
* 设备离线
|
||||
*
|
||||
* @param deviceUniqueId 设备唯一标识
|
||||
*/
|
||||
void deviceOffline(String deviceUniqueId);
|
||||
|
||||
/**
|
||||
* 更新设备最后消息时间
|
||||
*
|
||||
* @param deviceUniqueId 设备唯一标识
|
||||
* @param messageTime 消息时间
|
||||
*/
|
||||
void updateLastMessageTime(String deviceUniqueId, LocalDateTime messageTime);
|
||||
|
||||
/**
|
||||
* 检查设备是否在线
|
||||
*
|
||||
* @param deviceUniqueId 设备唯一标识
|
||||
* @return 是否在线
|
||||
*/
|
||||
boolean isDeviceOnline(String deviceUniqueId);
|
||||
|
||||
/**
|
||||
* 获取设备状态
|
||||
*
|
||||
* @param deviceUniqueId 设备唯一标识
|
||||
* @return 设备状态
|
||||
*/
|
||||
DeviceStateEnum getDeviceState(String deviceUniqueId);
|
||||
|
||||
/**
|
||||
* 获取在线设备列表
|
||||
*
|
||||
* @return 在线设备列表
|
||||
*/
|
||||
List<DeviceInfoDO> getOnlineDevices();
|
||||
|
||||
/**
|
||||
* 获取离线设备列表
|
||||
*
|
||||
* @return 离线设备列表
|
||||
*/
|
||||
List<DeviceInfoDO> getOfflineDevices();
|
||||
|
||||
/**
|
||||
* 获取指定状态的设备列表
|
||||
*
|
||||
* @param state 设备状态
|
||||
* @return 设备列表
|
||||
*/
|
||||
List<DeviceInfoDO> getDevicesByState(Integer state);
|
||||
|
||||
/**
|
||||
* 获取超时未上报的设备列表
|
||||
*
|
||||
* @param timeoutMinutes 超时分钟数
|
||||
* @return 超时设备列表
|
||||
*/
|
||||
List<DeviceInfoDO> getTimeoutDevices(int timeoutMinutes);
|
||||
|
||||
/**
|
||||
* 批量检查设备离线状态
|
||||
*
|
||||
* @param timeoutMinutes 离线超时分钟数
|
||||
* @return 检查到的离线设备数量
|
||||
*/
|
||||
int checkAndUpdateOfflineDevices(int timeoutMinutes);
|
||||
}
|
||||
@@ -0,0 +1,152 @@
|
||||
package cn.iocoder.yudao.module.farm.service.devicestatus;
|
||||
|
||||
import cn.iocoder.yudao.module.farm.dal.dataobject.deviceinfo.DeviceInfoDO;
|
||||
import cn.iocoder.yudao.module.farm.dal.mysql.deviceinfo.DeviceInfoMapper;
|
||||
import cn.iocoder.yudao.module.farm.enums.DeviceStateEnum;
|
||||
import jakarta.annotation.Resource;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.transaction.annotation.Transactional;
|
||||
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* 设备状态管理 Service 实现类
|
||||
*
|
||||
* @author 芋道源码
|
||||
*/
|
||||
@Service
|
||||
@Slf4j
|
||||
public class DeviceStatusServiceImpl implements DeviceStatusService {
|
||||
|
||||
@Resource
|
||||
private DeviceInfoMapper deviceInfoMapper;
|
||||
|
||||
@Override
|
||||
@Transactional
|
||||
public void updateDeviceState(String deviceUniqueId, Integer state) {
|
||||
updateDeviceState(deviceUniqueId, state, LocalDateTime.now());
|
||||
}
|
||||
|
||||
@Override
|
||||
@Transactional
|
||||
public void updateDeviceState(String deviceUniqueId, Integer state, LocalDateTime messageTime) {
|
||||
DeviceInfoDO device = deviceInfoMapper.selectOne(DeviceInfoDO::getDeviceUniqueId, deviceUniqueId);
|
||||
if (device == null) {
|
||||
log.warn("设备不存在:{}", deviceUniqueId);
|
||||
return;
|
||||
}
|
||||
|
||||
// 构建更新对象
|
||||
DeviceInfoDO updateObj = new DeviceInfoDO();
|
||||
updateObj.setId(device.getId());
|
||||
updateObj.setState(state);
|
||||
updateObj.setLastMessageTime(messageTime);
|
||||
|
||||
// 根据状态设置相应的时间字段
|
||||
if (DeviceStateEnum.ONLINE.getState().equals(state)) {
|
||||
// 设备上线
|
||||
updateObj.setOnlineTime(messageTime);
|
||||
// 如果是首次激活,设置激活时间
|
||||
if (device.getActiveTime() == null) {
|
||||
updateObj.setActiveTime(messageTime);
|
||||
}
|
||||
} else if (DeviceStateEnum.OFFLINE.getState().equals(state)) {
|
||||
// 设备离线
|
||||
updateObj.setOfflineTime(messageTime);
|
||||
}
|
||||
|
||||
deviceInfoMapper.updateById(updateObj);
|
||||
log.info("设备状态已更新:{} -> {}", deviceUniqueId, DeviceStateEnum.valueOf(state).getName());
|
||||
}
|
||||
|
||||
@Override
|
||||
public void deviceOnline(String deviceUniqueId) {
|
||||
updateDeviceState(deviceUniqueId, DeviceStateEnum.ONLINE.getState());
|
||||
}
|
||||
|
||||
@Override
|
||||
public void deviceOffline(String deviceUniqueId) {
|
||||
updateDeviceState(deviceUniqueId, DeviceStateEnum.OFFLINE.getState());
|
||||
}
|
||||
|
||||
@Override
|
||||
public void updateLastMessageTime(String deviceUniqueId, LocalDateTime messageTime) {
|
||||
DeviceInfoDO device = deviceInfoMapper.selectOne(DeviceInfoDO::getDeviceUniqueId, deviceUniqueId);
|
||||
if (device == null) {
|
||||
log.warn("设备不存在:{}", deviceUniqueId);
|
||||
return;
|
||||
}
|
||||
|
||||
DeviceInfoDO updateObj = new DeviceInfoDO();
|
||||
updateObj.setId(device.getId());
|
||||
updateObj.setLastMessageTime(messageTime);
|
||||
|
||||
// 如果设备当前是离线状态,收到消息后自动上线
|
||||
if (DeviceStateEnum.isOffline(device.getState()) || device.getState() == null) {
|
||||
updateObj.setState(DeviceStateEnum.ONLINE.getState());
|
||||
updateObj.setOnlineTime(messageTime);
|
||||
// 如果是首次激活,设置激活时间
|
||||
if (device.getActiveTime() == null) {
|
||||
updateObj.setActiveTime(messageTime);
|
||||
}
|
||||
}
|
||||
|
||||
deviceInfoMapper.updateById(updateObj);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean isDeviceOnline(String deviceUniqueId) {
|
||||
DeviceInfoDO device = deviceInfoMapper.selectOne(DeviceInfoDO::getDeviceUniqueId, deviceUniqueId);
|
||||
return device != null && DeviceStateEnum.isOnline(device.getState());
|
||||
}
|
||||
|
||||
@Override
|
||||
public DeviceStateEnum getDeviceState(String deviceUniqueId) {
|
||||
DeviceInfoDO device = deviceInfoMapper.selectOne(DeviceInfoDO::getDeviceUniqueId, deviceUniqueId);
|
||||
if (device == null || device.getState() == null) {
|
||||
return DeviceStateEnum.INACTIVE;
|
||||
}
|
||||
return DeviceStateEnum.valueOf(device.getState());
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<DeviceInfoDO> getOnlineDevices() {
|
||||
return getDevicesByState(DeviceStateEnum.ONLINE.getState());
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<DeviceInfoDO> getOfflineDevices() {
|
||||
return getDevicesByState(DeviceStateEnum.OFFLINE.getState());
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<DeviceInfoDO> getDevicesByState(Integer state) {
|
||||
return deviceInfoMapper.selectListByState(state);
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<DeviceInfoDO> getTimeoutDevices(int timeoutMinutes) {
|
||||
LocalDateTime timeoutTime = LocalDateTime.now().minusMinutes(timeoutMinutes);
|
||||
return deviceInfoMapper.selectTimeoutDevices(timeoutTime);
|
||||
}
|
||||
|
||||
@Override
|
||||
@Transactional
|
||||
public int checkAndUpdateOfflineDevices(int timeoutMinutes) {
|
||||
List<DeviceInfoDO> timeoutDevices = getTimeoutDevices(timeoutMinutes);
|
||||
int count = 0;
|
||||
|
||||
for (DeviceInfoDO device : timeoutDevices) {
|
||||
// 只处理当前在线的设备
|
||||
if (DeviceStateEnum.isOnline(device.getState())) {
|
||||
deviceOffline(device.getDeviceUniqueId());
|
||||
count++;
|
||||
log.info("设备超时离线:{}", device.getDeviceUniqueId());
|
||||
}
|
||||
}
|
||||
|
||||
return count;
|
||||
}
|
||||
}
|
||||
@@ -1,7 +1,10 @@
|
||||
package cn.iocoder.yudao.module.farm.service.mqtt;
|
||||
|
||||
import cn.iocoder.yudao.module.farm.dal.dataobject.deviceinfo.DeviceInfoDO;
|
||||
import cn.iocoder.yudao.module.farm.dal.dataobject.mqtt.MqttConnectionDO;
|
||||
import cn.iocoder.yudao.module.farm.dal.mysql.deviceinfo.DeviceInfoMapper;
|
||||
import cn.iocoder.yudao.module.farm.dal.mysql.mqtt.MqttConnectionMapper;
|
||||
import cn.iocoder.yudao.module.farm.service.devicestatus.DeviceStatusService;
|
||||
import cn.iocoder.yudao.module.farm.service.telemetry.DeviceTelemetryService;
|
||||
import jakarta.annotation.Resource;
|
||||
import org.eclipse.paho.client.mqttv3.IMqttMessageListener;
|
||||
@@ -19,19 +22,41 @@ public class DeviceMqttService {
|
||||
@Resource
|
||||
private DeviceInfoMapper deviceInfoMapper;
|
||||
@Resource
|
||||
private MqttConnectionMapper mqttConnectionMapper;
|
||||
@Resource
|
||||
private MqttConnectionService mqttConnectionService;
|
||||
@Resource
|
||||
private DeviceTelemetryService deviceTelemetryService;
|
||||
@Resource
|
||||
private DeviceStatusService deviceStatusService;
|
||||
|
||||
private final Set<String> subscribed = ConcurrentHashMap.newKeySet();
|
||||
|
||||
// 上报主题
|
||||
private String telemetryTopic(String deviceUniqueId) {
|
||||
DeviceInfoDO device = getDevice(deviceUniqueId);
|
||||
if (device != null) {
|
||||
String tenantId = String.valueOf(device.getTenantId());
|
||||
MqttConnectionDO connection = mqttConnectionMapper.selectByDeviceIdAndTenant(deviceUniqueId, tenantId);
|
||||
if (connection != null && connection.getTelemetryTopic() != null) {
|
||||
return connection.getTelemetryTopic().replace("{deviceId}", deviceUniqueId);
|
||||
}
|
||||
}
|
||||
// 默认主题格式
|
||||
return "farm/" + deviceUniqueId + "/telemetry";
|
||||
}
|
||||
|
||||
// 下发指令
|
||||
private String commandTopic(String deviceUniqueId) {
|
||||
DeviceInfoDO device = getDevice(deviceUniqueId);
|
||||
if (device != null) {
|
||||
String tenantId = String.valueOf(device.getTenantId());
|
||||
MqttConnectionDO connection = mqttConnectionMapper.selectByDeviceIdAndTenant(deviceUniqueId, tenantId);
|
||||
if (connection != null && connection.getCommandTopic() != null) {
|
||||
return connection.getCommandTopic().replace("{deviceId}", deviceUniqueId);
|
||||
}
|
||||
}
|
||||
// 默认主题格式
|
||||
return "farm/" + deviceUniqueId + "/cmd";
|
||||
}
|
||||
|
||||
@@ -71,8 +96,32 @@ public class DeviceMqttService {
|
||||
public void subscribeTelemetryAndPersist(String deviceUniqueId) throws MqttException {
|
||||
this.subscribeTelemetry(deviceUniqueId, (topic, message) -> {
|
||||
String payload = new String(message.getPayload(), StandardCharsets.UTF_8);
|
||||
// 保存遥测数据
|
||||
deviceTelemetryService.saveTelemetry(deviceUniqueId, payload);
|
||||
// 更新设备最后消息时间(会自动处理设备上线逻辑)
|
||||
deviceStatusService.updateLastMessageTime(deviceUniqueId, java.time.LocalDateTime.now());
|
||||
});
|
||||
|
||||
// 同时订阅LWT主题,处理设备异常离线
|
||||
subscribeLWT(deviceUniqueId);
|
||||
}
|
||||
|
||||
// 订阅LWT主题,处理设备离线
|
||||
public void subscribeLWT(String deviceUniqueId) throws MqttException {
|
||||
MqttClient client = mqttConnectionService.getClientByDeviceId(deviceUniqueId);
|
||||
if (client == null) return;
|
||||
|
||||
String lwtTopic = "farm/" + deviceUniqueId + "/lwt";
|
||||
String key = client.getClientId() + "|" + lwtTopic;
|
||||
|
||||
if (subscribed.add(key)) {
|
||||
client.subscribe(lwtTopic, 1, (topic, message) -> {
|
||||
String payload = new String(message.getPayload(), StandardCharsets.UTF_8);
|
||||
// LWT消息表示设备异常离线
|
||||
deviceStatusService.deviceOffline(deviceUniqueId);
|
||||
System.out.println("设备异常离线:" + deviceUniqueId + ",LWT消息:" + payload);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
public DeviceInfoDO getDevice(String deviceUniqueId) {
|
||||
|
||||
@@ -0,0 +1,85 @@
|
||||
package cn.iocoder.yudao.module.farm.service.mqtt;
|
||||
|
||||
import cn.iocoder.yudao.module.farm.controller.admin.mqtt.vo.MqttConnectionPageReqVO;
|
||||
import cn.iocoder.yudao.module.farm.controller.admin.mqtt.vo.MqttConnectionSaveReqVO;
|
||||
import cn.iocoder.yudao.module.farm.dal.dataobject.mqtt.MqttConnectionDO;
|
||||
import cn.iocoder.yudao.framework.common.pojo.PageResult;
|
||||
import jakarta.validation.Valid;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* MQTT连接管理 Service 接口
|
||||
*
|
||||
* @author 芋道源码
|
||||
*/
|
||||
public interface MqttConnectionManagementService {
|
||||
|
||||
/**
|
||||
* 创建MQTT连接配置
|
||||
*
|
||||
* @param createReqVO 创建信息
|
||||
* @return 编号
|
||||
*/
|
||||
Long createMqttConnection(@Valid MqttConnectionSaveReqVO createReqVO);
|
||||
|
||||
/**
|
||||
* 更新MQTT连接配置
|
||||
*
|
||||
* @param updateReqVO 更新信息
|
||||
*/
|
||||
void updateMqttConnection(@Valid MqttConnectionSaveReqVO updateReqVO);
|
||||
|
||||
/**
|
||||
* 删除MQTT连接配置
|
||||
*
|
||||
* @param id 编号
|
||||
*/
|
||||
void deleteMqttConnection(Long id);
|
||||
|
||||
/**
|
||||
* 批量删除MQTT连接配置
|
||||
*
|
||||
* @param ids 编号列表
|
||||
*/
|
||||
void deleteMqttConnectionListByIds(List<Long> ids);
|
||||
|
||||
/**
|
||||
* 获得MQTT连接配置
|
||||
*
|
||||
* @param id 编号
|
||||
* @return MQTT连接配置
|
||||
*/
|
||||
MqttConnectionDO getMqttConnection(Long id);
|
||||
|
||||
/**
|
||||
* 获得MQTT连接配置分页
|
||||
*
|
||||
* @param pageReqVO 分页查询
|
||||
* @return MQTT连接配置分页
|
||||
*/
|
||||
PageResult<MqttConnectionDO> getMqttConnectionPage(MqttConnectionPageReqVO pageReqVO);
|
||||
|
||||
/**
|
||||
* 根据设备ID获取MQTT连接配置
|
||||
*
|
||||
* @param deviceId 设备ID
|
||||
* @return MQTT连接配置
|
||||
*/
|
||||
MqttConnectionDO getMqttConnectionByDeviceId(String deviceId);
|
||||
|
||||
/**
|
||||
* 测试MQTT连接
|
||||
*
|
||||
* @param id 连接配置ID
|
||||
* @return 连接结果
|
||||
*/
|
||||
boolean testMqttConnection(Long id);
|
||||
|
||||
/**
|
||||
* 重新加载MQTT连接
|
||||
*
|
||||
* @param deviceId 设备ID
|
||||
*/
|
||||
void reloadMqttConnection(String deviceId);
|
||||
}
|
||||
@@ -0,0 +1,141 @@
|
||||
package cn.iocoder.yudao.module.farm.service.mqtt;
|
||||
|
||||
import cn.iocoder.yudao.module.farm.controller.admin.mqtt.vo.MqttConnectionPageReqVO;
|
||||
import cn.iocoder.yudao.module.farm.controller.admin.mqtt.vo.MqttConnectionSaveReqVO;
|
||||
import cn.iocoder.yudao.module.farm.dal.dataobject.mqtt.MqttConnectionDO;
|
||||
import cn.iocoder.yudao.module.farm.dal.mysql.mqtt.MqttConnectionMapper;
|
||||
import cn.iocoder.yudao.framework.common.pojo.PageResult;
|
||||
import cn.iocoder.yudao.framework.common.util.object.BeanUtils;
|
||||
import cn.iocoder.yudao.framework.security.core.LoginUser;
|
||||
import jakarta.annotation.Resource;
|
||||
import org.eclipse.paho.client.mqttv3.MqttClient;
|
||||
import org.eclipse.paho.client.mqttv3.MqttException;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.validation.annotation.Validated;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
import static cn.iocoder.yudao.framework.common.exception.util.ServiceExceptionUtil.exception;
|
||||
import static cn.iocoder.yudao.framework.security.core.util.SecurityFrameworkUtils.getLoginUser;
|
||||
import static cn.iocoder.yudao.module.farm.enums.ErrorCodeConstants.MQTT_CONNECTION_NOT_EXISTS;
|
||||
|
||||
/**
|
||||
* MQTT连接管理 Service 实现类
|
||||
*
|
||||
* @author 芋道源码
|
||||
*/
|
||||
@Service
|
||||
@Validated
|
||||
public class MqttConnectionManagementServiceImpl implements MqttConnectionManagementService {
|
||||
|
||||
@Resource
|
||||
private MqttConnectionMapper mqttConnectionMapper;
|
||||
|
||||
@Resource
|
||||
private MqttConnectionService mqttConnectionService;
|
||||
|
||||
@Override
|
||||
public Long createMqttConnection(MqttConnectionSaveReqVO createReqVO) {
|
||||
LoginUser user = getLoginUser();
|
||||
if (user == null) {
|
||||
throw new IllegalStateException("用户未登录");
|
||||
}
|
||||
// 插入
|
||||
MqttConnectionDO mqttConnection = BeanUtils.toBean(createReqVO, MqttConnectionDO.class);
|
||||
mqttConnection.setTenantId(String.valueOf(user.getTenantId()));
|
||||
mqttConnectionMapper.insert(mqttConnection);
|
||||
// 返回
|
||||
return mqttConnection.getId();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void updateMqttConnection(MqttConnectionSaveReqVO updateReqVO) {
|
||||
// 校验存在
|
||||
validateMqttConnectionExists(updateReqVO.getId());
|
||||
// 更新
|
||||
MqttConnectionDO updateObj = BeanUtils.toBean(updateReqVO, MqttConnectionDO.class);
|
||||
mqttConnectionMapper.updateById(updateObj);
|
||||
|
||||
// 重新加载连接
|
||||
MqttConnectionDO connection = mqttConnectionMapper.selectById(updateReqVO.getId());
|
||||
if (connection != null && connection.getDeviceId() != null) {
|
||||
mqttConnectionService.reload(connection.getDeviceId());
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void deleteMqttConnection(Long id) {
|
||||
// 校验存在
|
||||
MqttConnectionDO connection = validateMqttConnectionExists(id);
|
||||
// 删除
|
||||
mqttConnectionMapper.deleteById(id);
|
||||
|
||||
// 清理连接
|
||||
if (connection.getDeviceId() != null) {
|
||||
mqttConnectionService.reload(connection.getDeviceId());
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void deleteMqttConnectionListByIds(List<Long> ids) {
|
||||
// 获取要删除的连接信息
|
||||
List<MqttConnectionDO> connections = mqttConnectionMapper.selectByIds(ids);
|
||||
// 删除
|
||||
mqttConnectionMapper.deleteByIds(ids);
|
||||
|
||||
// 清理连接
|
||||
connections.forEach(connection -> {
|
||||
if (connection.getDeviceId() != null) {
|
||||
mqttConnectionService.reload(connection.getDeviceId());
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
private MqttConnectionDO validateMqttConnectionExists(Long id) {
|
||||
MqttConnectionDO connection = mqttConnectionMapper.selectById(id);
|
||||
if (connection == null) {
|
||||
throw exception(MQTT_CONNECTION_NOT_EXISTS);
|
||||
}
|
||||
return connection;
|
||||
}
|
||||
|
||||
@Override
|
||||
public MqttConnectionDO getMqttConnection(Long id) {
|
||||
return mqttConnectionMapper.selectById(id);
|
||||
}
|
||||
|
||||
@Override
|
||||
public PageResult<MqttConnectionDO> getMqttConnectionPage(MqttConnectionPageReqVO pageReqVO) {
|
||||
return mqttConnectionMapper.selectPage(pageReqVO);
|
||||
}
|
||||
|
||||
@Override
|
||||
public MqttConnectionDO getMqttConnectionByDeviceId(String deviceId) {
|
||||
LoginUser user = getLoginUser();
|
||||
if (user == null) {
|
||||
return null;
|
||||
}
|
||||
String tenantId = String.valueOf(user.getTenantId());
|
||||
return mqttConnectionMapper.selectByDeviceIdAndTenant(deviceId, tenantId);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean testMqttConnection(Long id) {
|
||||
MqttConnectionDO connection = getMqttConnection(id);
|
||||
if (connection == null || connection.getDeviceId() == null) {
|
||||
return false;
|
||||
}
|
||||
|
||||
try {
|
||||
MqttClient client = mqttConnectionService.getClientByDeviceId(connection.getDeviceId(), connection.getTenantId());
|
||||
return client != null && client.isConnected();
|
||||
} catch (MqttException e) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void reloadMqttConnection(String deviceId) {
|
||||
mqttConnectionService.reload(deviceId);
|
||||
}
|
||||
}
|
||||
@@ -56,8 +56,10 @@ public class MqttConnectionService {
|
||||
options.setPassword(cfg.getPassword() != null ? cfg.getPassword().toCharArray() : new char[0]);
|
||||
}
|
||||
try {
|
||||
// 设置 Last Will Testament (LWT) - 设备异常断开时的遗言消息
|
||||
String lwtTopic = "farm/" + deviceId + "/lwt";
|
||||
MqttMessage lwt = new MqttMessage("offline".getBytes(StandardCharsets.UTF_8));
|
||||
String lwtMessage = "{\"deviceId\":\"" + deviceId + "\",\"status\":\"offline\",\"timestamp\":" + System.currentTimeMillis() + "}";
|
||||
MqttMessage lwt = new MqttMessage(lwtMessage.getBytes(StandardCharsets.UTF_8));
|
||||
lwt.setQos(cfg.getQos() != null ? cfg.getQos() : 1);
|
||||
lwt.setRetained(true);
|
||||
options.setWill(lwtTopic, lwt.getPayload(), lwt.getQos(), lwt.isRetained());
|
||||
|
||||
@@ -0,0 +1,106 @@
|
||||
package cn.iocoder.yudao.module.farm.service.mqtt;
|
||||
|
||||
import cn.iocoder.yudao.module.farm.dal.dataobject.mqtt.MqttConnectionDO;
|
||||
import cn.iocoder.yudao.module.farm.dal.mysql.mqtt.MqttConnectionMapper;
|
||||
import jakarta.annotation.Resource;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.boot.CommandLineRunner;
|
||||
import org.springframework.stereotype.Service;
|
||||
|
||||
/**
|
||||
* MQTT设备初始化服务
|
||||
* 用于初始化设备连接配置
|
||||
*
|
||||
* @author 芋道源码
|
||||
*/
|
||||
@Service
|
||||
@Slf4j
|
||||
public class MqttDeviceInitService implements CommandLineRunner {
|
||||
|
||||
@Resource
|
||||
private MqttConnectionMapper mqttConnectionMapper;
|
||||
|
||||
@Override
|
||||
public void run(String... args) throws Exception {
|
||||
initDevice2409AA01();
|
||||
}
|
||||
|
||||
/**
|
||||
* 初始化设备2409AA01的MQTT连接配置
|
||||
*/
|
||||
private void initDevice2409AA01() {
|
||||
String deviceId = "2409AA01";
|
||||
String tenantId = "1"; // 默认租户ID,实际使用时需要根据具体情况调整
|
||||
|
||||
// 检查是否已经存在配置
|
||||
MqttConnectionDO existing = mqttConnectionMapper.selectByDeviceIdAndTenant(deviceId, tenantId);
|
||||
if (existing != null) {
|
||||
log.info("设备 {} 的MQTT连接配置已存在,跳过初始化", deviceId);
|
||||
return;
|
||||
}
|
||||
|
||||
// 创建新的连接配置
|
||||
MqttConnectionDO connection = MqttConnectionDO.builder()
|
||||
.tenantId(tenantId)
|
||||
.name("设备连接-" + deviceId)
|
||||
.deviceId(deviceId)
|
||||
.broker("tcp://mqtt.api.mimirii.com:1883")
|
||||
.clientId(deviceId)
|
||||
.username("mimir_user")
|
||||
.password("mimir123456@")
|
||||
.cleanSession(Boolean.TRUE)
|
||||
.keepAlive(30)
|
||||
.qos(1)
|
||||
.reconnect(Boolean.TRUE)
|
||||
.status(Boolean.TRUE)
|
||||
.telemetryTopic("/WFM/{deviceId}/up")
|
||||
.commandTopic("/WFM/{deviceId}/down")
|
||||
.build();
|
||||
|
||||
mqttConnectionMapper.insert(connection);
|
||||
log.info("已成功初始化设备 {} 的MQTT连接配置", deviceId);
|
||||
}
|
||||
|
||||
/**
|
||||
* 手动初始化设备连接配置(供API调用)
|
||||
*
|
||||
* @param deviceId 设备ID
|
||||
* @param tenantId 租户ID
|
||||
* @param broker MQTT服务器地址
|
||||
* @param username 用户名
|
||||
* @param password 密码
|
||||
* @param telemetryTopic 上报主题模板
|
||||
* @param commandTopic 下发主题模板
|
||||
*/
|
||||
public void initDeviceConnection(String deviceId, String tenantId, String broker,
|
||||
String username, String password,
|
||||
String telemetryTopic, String commandTopic) {
|
||||
// 检查是否已经存在配置
|
||||
MqttConnectionDO existing = mqttConnectionMapper.selectByDeviceIdAndTenant(deviceId, tenantId);
|
||||
if (existing != null) {
|
||||
log.warn("设备 {} 的MQTT连接配置已存在", deviceId);
|
||||
return;
|
||||
}
|
||||
|
||||
// 创建新的连接配置
|
||||
MqttConnectionDO connection = MqttConnectionDO.builder()
|
||||
.tenantId(tenantId)
|
||||
.name("设备连接-" + deviceId)
|
||||
.deviceId(deviceId)
|
||||
.broker(broker)
|
||||
.clientId(deviceId)
|
||||
.username(username)
|
||||
.password(password)
|
||||
.cleanSession(Boolean.TRUE)
|
||||
.keepAlive(30)
|
||||
.qos(1)
|
||||
.reconnect(Boolean.TRUE)
|
||||
.status(Boolean.TRUE)
|
||||
.telemetryTopic(telemetryTopic)
|
||||
.commandTopic(commandTopic)
|
||||
.build();
|
||||
|
||||
mqttConnectionMapper.insert(connection);
|
||||
log.info("已成功初始化设备 {} 的MQTT连接配置", deviceId);
|
||||
}
|
||||
}
|
||||
@@ -283,6 +283,7 @@ yudao:
|
||||
- /admin-api/app/traceability-code/scan
|
||||
ignore-tables:
|
||||
- region
|
||||
- farm_mqtt_connection
|
||||
- farm_production_plan
|
||||
- farm_species_classify_library
|
||||
ignore-caches:
|
||||
|
||||
Reference in New Issue
Block a user