diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/deviceinfo/DeviceInfoController.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/deviceinfo/DeviceInfoController.java index 1ea0306f..93623b0d 100644 --- a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/deviceinfo/DeviceInfoController.java +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/deviceinfo/DeviceInfoController.java @@ -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 = "创建设备信息") @@ -90,12 +89,23 @@ public class DeviceInfoController { @Operation(summary = "导出设备信息 Excel") @ApiAccessLog(operateType = EXPORT) public void exportDeviceInfoExcel(@Valid DeviceInfoPageReqVO pageReqVO, - HttpServletResponse response) throws IOException { + HttpServletResponse response) throws IOException { pageReqVO.setPageSize(PageParam.PAGE_SIZE_NONE); List list = deviceInfoService.getDeviceInfoPage(pageReqVO).getList(); // 导出 Excel ExcelUtils.write(response, "设备信息.xls", "数据", DeviceInfoRespVO.class, - BeanUtils.toBean(list, DeviceInfoRespVO.class)); + BeanUtils.toBean(list, DeviceInfoRespVO.class)); } + @PostMapping("/startSub") + @Operation(summary = "订阅设备上报并持久化") + public CommonResult startSub(@Valid @RequestBody MqttConnectionDO createReqVO) { + try { + deviceMqttService.subscribeTelemetryAndPersist(createReqVO.getDeviceId()); + return success(true); + } catch (MqttException e) { + return error(500, "设备连接失败"); + } + + } } \ No newline at end of file diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/devicestatus/DeviceStatusController.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/devicestatus/DeviceStatusController.java new file mode 100644 index 00000000..b21ea3bc --- /dev/null +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/devicestatus/DeviceStatusController.java @@ -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 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 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> getOnlineDevices() { + List devices = deviceStatusService.getOnlineDevices(); + List 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> getOfflineDevices() { + List devices = deviceStatusService.getOfflineDevices(); + List 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> getTimeoutDevices( + @RequestParam(value = "timeoutMinutes", defaultValue = "10") int timeoutMinutes) { + List devices = deviceStatusService.getTimeoutDevices(timeoutMinutes); + List 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 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 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 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 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 checkOfflineDevices( + @RequestParam(value = "timeoutMinutes", defaultValue = "10") int timeoutMinutes) { + int count = deviceStatusService.checkAndUpdateOfflineDevices(timeoutMinutes); + return success(count); + } +} diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/devicestatus/vo/DeviceStatusPageReqVO.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/devicestatus/vo/DeviceStatusPageReqVO.java new file mode 100644 index 00000000..0133a788 --- /dev/null +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/devicestatus/vo/DeviceStatusPageReqVO.java @@ -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; +} diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/devicestatus/vo/DeviceStatusRespVO.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/devicestatus/vo/DeviceStatusRespVO.java new file mode 100644 index 00000000..bfcca561 --- /dev/null +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/devicestatus/vo/DeviceStatusRespVO.java @@ -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; +} diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/devicestatus/vo/DeviceStatusStatisticsRespVO.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/devicestatus/vo/DeviceStatusStatisticsRespVO.java new file mode 100644 index 00000000..e8faba29 --- /dev/null +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/devicestatus/vo/DeviceStatusStatisticsRespVO.java @@ -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; +} diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/mqtt/MqttConnectionController.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/mqtt/MqttConnectionController.java new file mode 100644 index 00000000..d6f6ddd2 --- /dev/null +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/mqtt/MqttConnectionController.java @@ -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 createMqttConnection(@Valid @RequestBody MqttConnectionSaveReqVO createReqVO) { + return success(mqttConnectionManagementService.createMqttConnection(createReqVO)); + } + + @PutMapping("/update") + @Operation(summary = "更新MQTT连接配置") + @PreAuthorize("@ss.hasPermission('farm:mqtt-connection:update')") + public CommonResult 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 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 deleteMqttConnectionListByIds(@RequestParam("ids") List 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 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> getMqttConnectionPage(@Valid MqttConnectionPageReqVO pageReqVO) { + PageResult 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 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 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 reloadMqttConnection(@PathVariable("deviceId") String deviceId) { + mqttConnectionManagementService.reloadMqttConnection(deviceId); + return success(true); + } +} diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/mqtt/MqttTestController.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/mqtt/MqttTestController.java new file mode 100644 index 00000000..0717c7b0 --- /dev/null +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/mqtt/MqttTestController.java @@ -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 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 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 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 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 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 reloadConnection(@PathVariable("deviceId") String deviceId) { + mqttConnectionService.reload(deviceId); + return success("连接已重新加载"); + } +} diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/mqtt/vo/MqttConnectionPageReqVO.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/mqtt/vo/MqttConnectionPageReqVO.java new file mode 100644 index 00000000..b9833522 --- /dev/null +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/mqtt/vo/MqttConnectionPageReqVO.java @@ -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; +} diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/mqtt/vo/MqttConnectionRespVO.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/mqtt/vo/MqttConnectionRespVO.java new file mode 100644 index 00000000..1643f580 --- /dev/null +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/mqtt/vo/MqttConnectionRespVO.java @@ -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; +} diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/mqtt/vo/MqttConnectionSaveReqVO.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/mqtt/vo/MqttConnectionSaveReqVO.java new file mode 100644 index 00000000..f50cb5e9 --- /dev/null +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/controller/admin/mqtt/vo/MqttConnectionSaveReqVO.java @@ -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; +} diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/dataobject/deviceinfo/DeviceInfoDO.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/dataobject/deviceinfo/DeviceInfoDO.java index 1e3ff71b..e78fb170 100644 --- a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/dataobject/deviceinfo/DeviceInfoDO.java +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/dataobject/deviceinfo/DeviceInfoDO.java @@ -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; + /** * 数据快照集合(业务层使用) */ diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/dataobject/modelenvironmentindex/ModelEnvironmentIndexDO.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/dataobject/modelenvironmentindex/ModelEnvironmentIndexDO.java index 4447e370..26ac2c0d 100644 --- a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/dataobject/modelenvironmentindex/ModelEnvironmentIndexDO.java +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/dataobject/modelenvironmentindex/ModelEnvironmentIndexDO.java @@ -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; /** * 环境指标 diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/dataobject/mqtt/MqttConnectionDO.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/dataobject/mqtt/MqttConnectionDO.java index cf907fe0..68bcd58e 100644 --- a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/dataobject/mqtt/MqttConnectionDO.java +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/dataobject/mqtt/MqttConnectionDO.java @@ -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; } diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/mysql/deviceinfo/DeviceInfoMapper.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/mysql/deviceinfo/DeviceInfoMapper.java index c7e96d05..ff96269c 100644 --- a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/mysql/deviceinfo/DeviceInfoMapper.java +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/mysql/deviceinfo/DeviceInfoMapper.java @@ -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 { .orderByDesc(DeviceInfoDO::getCreateTime)); } + /** + * 根据设备状态查询设备列表 + */ + default List selectListByState(Integer state) { + return selectList(new LambdaQueryWrapperX() + .eq(DeviceInfoDO::getState, state)); + } + + /** + * 查询超时未上报的设备列表 + */ + default List selectTimeoutDevices(@Param("timeoutTime") LocalDateTime timeoutTime) { + return selectList(new LambdaQueryWrapperX() + .eq(DeviceInfoDO::getState, 1) // 只查询在线状态的设备 + .and(wrapper -> wrapper + .isNull(DeviceInfoDO::getLastMessageTime) // 从未上报过消息 + .or() + .lt(DeviceInfoDO::getLastMessageTime, timeoutTime) // 或者最后消息时间超时 + )); + } + + /** + * 统计各状态设备数量 + */ + default Long selectCountByState(Integer state) { + return selectCount(new LambdaQueryWrapperX() + .eq(DeviceInfoDO::getState, state)); + } + } \ No newline at end of file diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/mysql/mqtt/MqttConnectionMapper.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/mysql/mqtt/MqttConnectionMapper.java index cea5038d..6cfde09a 100644 --- a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/mysql/mqtt/MqttConnectionMapper.java +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/mysql/mqtt/MqttConnectionMapper.java @@ -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 { .eq(MqttConnectionDO::getStatus, true) ); } + + default PageResult selectPage(MqttConnectionPageReqVO reqVO) { + return selectPage(reqVO, new LambdaQueryWrapperX() + .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)); + } } diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/enums/DeviceStateEnum.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/enums/DeviceStateEnum.java new file mode 100644 index 00000000..12db5fb7 --- /dev/null +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/enums/DeviceStateEnum.java @@ -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 { + + 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); + } +} diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/enums/ErrorCodeConstants.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/enums/ErrorCodeConstants.java index 296ffc90..06d654bc 100644 --- a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/enums/ErrorCodeConstants.java +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/enums/ErrorCodeConstants.java @@ -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连接配置不存在"); } \ No newline at end of file diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/job/DeviceOfflineCheckJob.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/job/DeviceOfflineCheckJob.java new file mode 100644 index 00000000..7a579a3c --- /dev/null +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/job/DeviceOfflineCheckJob.java @@ -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 + *

+ * 检测逻辑:设备最后消息时间超过一定时间,则认为设备离线 + * + * @author 芋道源码 + */ +@Component +@Slf4j +public class DeviceOfflineCheckJob implements JobHandler { + + /** + * 设备离线超时时间(分钟) + *

+ * 可以通过定时任务参数配置,默认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 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 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; + } +} diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/devicestatus/DeviceStatusService.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/devicestatus/DeviceStatusService.java new file mode 100644 index 00000000..cfc701f2 --- /dev/null +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/devicestatus/DeviceStatusService.java @@ -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 getOnlineDevices(); + + /** + * 获取离线设备列表 + * + * @return 离线设备列表 + */ + List getOfflineDevices(); + + /** + * 获取指定状态的设备列表 + * + * @param state 设备状态 + * @return 设备列表 + */ + List getDevicesByState(Integer state); + + /** + * 获取超时未上报的设备列表 + * + * @param timeoutMinutes 超时分钟数 + * @return 超时设备列表 + */ + List getTimeoutDevices(int timeoutMinutes); + + /** + * 批量检查设备离线状态 + * + * @param timeoutMinutes 离线超时分钟数 + * @return 检查到的离线设备数量 + */ + int checkAndUpdateOfflineDevices(int timeoutMinutes); +} diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/devicestatus/DeviceStatusServiceImpl.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/devicestatus/DeviceStatusServiceImpl.java new file mode 100644 index 00000000..3481a431 --- /dev/null +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/devicestatus/DeviceStatusServiceImpl.java @@ -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 getOnlineDevices() { + return getDevicesByState(DeviceStateEnum.ONLINE.getState()); + } + + @Override + public List getOfflineDevices() { + return getDevicesByState(DeviceStateEnum.OFFLINE.getState()); + } + + @Override + public List getDevicesByState(Integer state) { + return deviceInfoMapper.selectListByState(state); + } + + @Override + public List getTimeoutDevices(int timeoutMinutes) { + LocalDateTime timeoutTime = LocalDateTime.now().minusMinutes(timeoutMinutes); + return deviceInfoMapper.selectTimeoutDevices(timeoutTime); + } + + @Override + @Transactional + public int checkAndUpdateOfflineDevices(int timeoutMinutes) { + List 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; + } +} diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/DeviceMqttService.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/DeviceMqttService.java index 24e4935e..af638137 100644 --- a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/DeviceMqttService.java +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/DeviceMqttService.java @@ -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 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) { diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttConnectionManagementService.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttConnectionManagementService.java new file mode 100644 index 00000000..48486ae2 --- /dev/null +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttConnectionManagementService.java @@ -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 ids); + + /** + * 获得MQTT连接配置 + * + * @param id 编号 + * @return MQTT连接配置 + */ + MqttConnectionDO getMqttConnection(Long id); + + /** + * 获得MQTT连接配置分页 + * + * @param pageReqVO 分页查询 + * @return MQTT连接配置分页 + */ + PageResult 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); +} diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttConnectionManagementServiceImpl.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttConnectionManagementServiceImpl.java new file mode 100644 index 00000000..1154207f --- /dev/null +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttConnectionManagementServiceImpl.java @@ -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 ids) { + // 获取要删除的连接信息 + List 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 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); + } +} diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttConnectionService.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttConnectionService.java index 8b80a79b..f785ee95 100644 --- a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttConnectionService.java +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttConnectionService.java @@ -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()); diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttDeviceInitService.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttDeviceInitService.java new file mode 100644 index 00000000..7d44c8af --- /dev/null +++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttDeviceInitService.java @@ -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); + } +} diff --git a/yudao-server/src/main/resources/application.yaml b/yudao-server/src/main/resources/application.yaml index 765b0521..e47c9d88 100644 --- a/yudao-server/src/main/resources/application.yaml +++ b/yudao-server/src/main/resources/application.yaml @@ -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: