Merge remote-tracking branch 'origin/main'

This commit is contained in:
达富斌
2025-08-25 16:04:49 +08:00
14 changed files with 1297 additions and 17 deletions

View File

@@ -0,0 +1,96 @@
package cn.iocoder.yudao.module.farm.controller.admin.websocket;
import cn.iocoder.yudao.framework.common.pojo.CommonResult;
import cn.iocoder.yudao.module.farm.service.messagepush.FarmMessagePushService;
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 jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.NotNull;
import static cn.iocoder.yudao.framework.common.pojo.CommonResult.success;
/**
* 农场WebSocket消息推送 Controller
*
* @author 芋道源码
*/
@Tag(name = "管理后台 - 农场WebSocket消息推送")
@RestController
@RequestMapping("/farm/websocket")
@Validated
@Slf4j
public class FarmWebSocketController {
@Resource
private FarmMessagePushService farmMessagePushService;
@PostMapping("/push-device-alarm")
@Operation(summary = "推送设备报警消息")
@PreAuthorize("@ss.hasPermission('farm:websocket:push')")
public CommonResult<String> pushDeviceAlarm(
@RequestParam @NotBlank(message = "设备ID不能为空") String deviceId,
@RequestParam @NotBlank(message = "设备名称不能为空") String deviceName,
@RequestParam @NotNull(message = "报警类型不能为空") Integer alarmType,
@RequestParam @NotBlank(message = "报警内容不能为空") String alarmContent,
@RequestParam(required = false) String baseId,
@RequestParam(required = false) String plotId) {
log.info("[pushDeviceAlarm][推送设备报警消息][deviceId: {}][alarmType: {}]", deviceId, alarmType);
farmMessagePushService.pushDeviceAlarm(deviceId, deviceName, alarmType, alarmContent, baseId, plotId);
return success("设备报警消息推送成功");
}
@PostMapping("/push-weather-warning")
@Operation(summary = "推送天气预警消息")
@PreAuthorize("@ss.hasPermission('farm:websocket:push')")
public CommonResult<String> pushWeatherWarning(
@RequestParam @NotBlank(message = "预警内容不能为空") String warningContent,
@RequestParam(required = false) String baseId,
@RequestParam(required = false) String plotId) {
log.info("[pushWeatherWarning][推送天气预警消息][baseId: {}][plotId: {}]", baseId, plotId);
farmMessagePushService.pushWeatherWarning(warningContent, baseId, plotId);
return success("天气预警消息推送成功");
}
@PostMapping("/push-system-notification")
@Operation(summary = "推送系统通知消息")
@PreAuthorize("@ss.hasPermission('farm:websocket:push')")
public CommonResult<String> pushSystemNotification(
@RequestParam @NotBlank(message = "标题不能为空") String title,
@RequestParam @NotBlank(message = "内容不能为空") String content,
@RequestParam(required = false) String baseId,
@RequestParam(required = false) String plotId) {
log.info("[pushSystemNotification][推送系统通知消息][title: {}][baseId: {}][plotId: {}]", title, baseId, plotId);
farmMessagePushService.pushSystemNotification(title, content, baseId, plotId);
return success("系统通知消息推送成功");
}
@PostMapping("/push-test-message")
@Operation(summary = "推送测试消息")
@PreAuthorize("@ss.hasPermission('farm:websocket:push')")
public CommonResult<String> pushTestMessage(
@RequestParam(defaultValue = "测试消息") String title,
@RequestParam(defaultValue = "这是一条测试消息,用于验证WebSocket连接是否正常。") String content) {
log.info("[pushTestMessage][推送测试消息][title: {}]", title);
farmMessagePushService.pushSystemNotification(title, content, null, null);
return success("测试消息推送成功");
}
}

View File

@@ -0,0 +1,100 @@
package cn.iocoder.yudao.module.farm.controller.admin.websocket.vo;
import lombok.Data;
import lombok.Builder;
import lombok.NoArgsConstructor;
import lombok.AllArgsConstructor;
import java.time.LocalDateTime;
/**
* 农场WebSocket消息 VO
*
* @author 芋道源码
*/
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class FarmWebSocketMessageVO {
/**
* 消息类型:1=报警,2=预警,3=消息
*/
private Integer messageType;
/**
* 消息类型名称
*/
private String messageTypeName;
/**
* 消息标题
*/
private String title;
/**
* 消息内容
*/
private String content;
/**
* 消息ID
*/
private String messageId;
/**
* 基地ID
*/
private String baseId;
/**
* 基地名称
*/
private String baseName;
/**
* 地块ID
*/
private String plotId;
/**
* 地块名称
*/
private String plotName;
/**
* 设备ID(如果是设备相关消息)
*/
private String deviceId;
/**
* 设备名称(如果是设备相关消息)
*/
private String deviceName;
/**
* 报警级别(如果是报警消息)
*/
private Integer alarmLevel;
/**
* 报警状态(如果是报警消息)
*/
private Integer alarmStatus;
/**
* 消息时间
*/
private LocalDateTime messageTime;
/**
* 是否已读
*/
private Boolean isRead;
/**
* 推送次数
*/
private Integer pushCount;
}

View File

@@ -0,0 +1,139 @@
package cn.iocoder.yudao.module.farm.dal.dataobject.devicealarm;
import com.baomidou.mybatisplus.annotation.*;
import cn.iocoder.yudao.framework.mybatis.core.dataobject.BaseDO;
import lombok.*;
import java.time.LocalDateTime;
/**
* 设备报警 DO
*
* @author 芋道源码
*/
@TableName("farm_device_alarm")
@KeySequence("farm_device_alarm_seq")
@Data
@EqualsAndHashCode(callSuper = true)
@ToString(callSuper = true)
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class DeviceAlarmDO extends BaseDO {
/**
* 报警ID
*/
@TableId(type = IdType.INPUT)
private String id;
/**
* 设备ID
*/
private String deviceId;
/**
* 设备名称
*/
private String deviceName;
/**
* 设备唯一标识
*/
private String deviceUniqueId;
/**
* 报警类型
*/
private Integer alarmType;
/**
* 报警类型名称
*/
private String alarmTypeName;
/**
* 报警级别:1=低,2=中,3=高,4=紧急
*/
private Integer alarmLevel;
/**
* 报警标题
*/
private String title;
/**
* 报警内容
*/
private String content;
/**
* 报警数据(JSON格式)
*/
private String alarmData;
/**
* 报警状态:0=未处理,1=已处理,2=已忽略
*/
private Integer status;
/**
* 处理时间
*/
private LocalDateTime processTime;
/**
* 处理人ID
*/
private String processUserId;
/**
* 处理人姓名
*/
private String processUserName;
/**
* 处理备注
*/
private String processRemark;
/**
* 基地ID
*/
private String baseId;
/**
* 基地名称
*/
private String baseName;
/**
* 地块ID
*/
private String plotId;
/**
* 地块名称
*/
private String plotName;
/**
* 租户ID
*/
private String tenantId;
/**
* 是否已推送:0=未推送,1=已推送
*/
private Integer pushStatus;
/**
* 推送时间
*/
private LocalDateTime pushTime;
/**
* 推送次数
*/
private Integer pushCount;
}

View File

@@ -0,0 +1,15 @@
package cn.iocoder.yudao.module.farm.dal.mysql.devicealarm;
import cn.iocoder.yudao.framework.mybatis.core.mapper.BaseMapperX;
import cn.iocoder.yudao.module.farm.dal.dataobject.devicealarm.DeviceAlarmDO;
import org.apache.ibatis.annotations.Mapper;
/**
* 设备报警 Mapper
*
* @author 芋道源码
*/
@Mapper
public interface DeviceAlarmMapper extends BaseMapperX<DeviceAlarmDO> {
}

View File

@@ -0,0 +1,40 @@
package cn.iocoder.yudao.module.farm.enums;
import lombok.AllArgsConstructor;
import lombok.Getter;
/**
* 报警状态枚举
*
* @author 芋道源码
*/
@Getter
@AllArgsConstructor
public enum AlarmStatusEnum {
UNRESOLVED(0, "未解除"),
RESOLVED(1, "已解除"),
NO_ALARM(2, "无报警");
/**
* 状态值
*/
private final Integer status;
/**
* 状态名称
*/
private final String name;
/**
* 根据状态值获取枚举
*/
public static AlarmStatusEnum getByStatus(Integer status) {
for (AlarmStatusEnum alarmStatus : values()) {
if (alarmStatus.getStatus().equals(status)) {
return alarmStatus;
}
}
return null;
}
}

View File

@@ -0,0 +1,53 @@
package cn.iocoder.yudao.module.farm.enums;
import lombok.AllArgsConstructor;
import lombok.Getter;
/**
* 设备报警类型枚举
*
* @author 芋道源码
*/
@Getter
@AllArgsConstructor
public enum DeviceAlarmTypeEnum {
OFFLINE(1, "设备离线"),
BATTERY_LOW(2, "电池电量低"),
SIGNAL_WEAK(3, "信号强度弱"),
TEMPERATURE_HIGH(4, "温度过高"),
TEMPERATURE_LOW(5, "温度过低"),
HUMIDITY_HIGH(6, "湿度过高"),
HUMIDITY_LOW(7, "湿度过低"),
SOIL_MOISTURE_LOW(8, "土壤湿度低"),
SOIL_MOISTURE_HIGH(9, "土壤湿度过高"),
PH_ABNORMAL(10, "pH值异常"),
NUTRIENT_LOW(11, "养分不足"),
PEST_DETECTED(12, "病虫害检测"),
DISEASE_DETECTED(13, "病害检测"),
WEATHER_ALERT(14, "天气预警"),
MAINTENANCE_DUE(15, "设备维护到期"),
CUSTOM(99, "自定义报警");
/**
* 类型值
*/
private final Integer type;
/**
* 类型名称
*/
private final String name;
/**
* 根据类型值获取枚举
*/
public static DeviceAlarmTypeEnum getByType(Integer type) {
for (DeviceAlarmTypeEnum alarmType : values()) {
if (alarmType.getType().equals(type)) {
return alarmType;
}
}
return null;
}
}

View File

@@ -107,6 +107,10 @@ public interface ErrorCodeConstants {
// ========== 模型管理 1_050_032_000 ==========
ErrorCode SCREEN_DISASTER_DATA_NOT_EXISTS = new ErrorCode(1_050_032_000, "灾害数据不存在");
// ========== 网格 1_050_033_000 ==========
ErrorCode GRID_INFO_NOT_EXISTS = new ErrorCode(1_050_033_000, "网格不存在");
// ========== 设备报警 1_050_033_000 ==========
ErrorCode DEVICE_ALARM_NOT_EXISTS = new ErrorCode(1_050_034_000, "设备报警不存在");
}

View File

@@ -0,0 +1,40 @@
package cn.iocoder.yudao.module.farm.enums;
import lombok.AllArgsConstructor;
import lombok.Getter;
/**
* 消息类型枚举
*
* @author 芋道源码
*/
@Getter
@AllArgsConstructor
public enum MessageTypeEnum {
ALARM(1, "报警"),
WARNING(2, "预警"),
MESSAGE(3, "消息");
/**
* 类型值
*/
private final Integer type;
/**
* 类型名称
*/
private final String name;
/**
* 根据类型值获取枚举
*/
public static MessageTypeEnum getByType(Integer type) {
for (MessageTypeEnum messageType : values()) {
if (messageType.getType().equals(type)) {
return messageType;
}
}
return null;
}
}

View File

@@ -0,0 +1,104 @@
package cn.iocoder.yudao.module.farm.service.devicealarm;
import cn.iocoder.yudao.module.farm.dal.dataobject.devicealarm.DeviceAlarmDO;
import cn.iocoder.yudao.framework.common.pojo.PageResult;
import java.util.List;
/**
* 设备报警 Service 接口
*
* @author 芋道源码
*/
public interface DeviceAlarmService {
/**
* 创建设备报警
*
* @param deviceAlarm 设备报警信息
* @return 报警ID
*/
String createDeviceAlarm(DeviceAlarmDO deviceAlarm);
/**
* 更新设备报警
*
* @param deviceAlarm 设备报警信息
*/
void updateDeviceAlarm(DeviceAlarmDO deviceAlarm);
/**
* 删除设备报警
*
* @param id 报警ID
*/
void deleteDeviceAlarm(String id);
/**
* 批量删除设备报警
*
* @param ids 报警ID列表
*/
void deleteDeviceAlarmListByIds(List<String> ids);
/**
* 获得设备报警
*
* @param id 报警ID
* @return 设备报警
*/
DeviceAlarmDO getDeviceAlarm(String id);
/**
* 获得设备报警分页
*
* @param pageReqVO 分页查询
* @return 设备报警分页
*/
PageResult<DeviceAlarmDO> getDeviceAlarmPage(Object pageReqVO);
/**
* 根据设备ID获取未处理的报警列表
*
* @param deviceId 设备ID
* @return 报警列表
*/
List<DeviceAlarmDO> getUnprocessedAlarmsByDeviceId(String deviceId);
/**
* 根据基地ID获取未处理的报警列表
*
* @param baseId 基地ID
* @return 报警列表
*/
List<DeviceAlarmDO> getUnprocessedAlarmsByBaseId(String baseId);
/**
* 处理设备报警
*
* @param id 报警ID
* @param processUserId 处理人ID
* @param processUserName 处理人姓名
* @param processRemark 处理备注
*/
void processDeviceAlarm(String id, String processUserId, String processUserName, String processRemark);
/**
* 忽略设备报警
*
* @param id 报警ID
* @param processUserId 处理人ID
* @param processUserName 处理人姓名
* @param processRemark 处理备注
*/
void ignoreDeviceAlarm(String id, String processUserId, String processUserName, String processRemark);
/**
* 检查设备是否需要报警
*
* @param deviceId 设备ID
* @param alarmType 报警类型
* @param alarmData 报警数据
*/
void checkAndCreateAlarm(String deviceId, Integer alarmType, String alarmData);
}

View File

@@ -0,0 +1,251 @@
package cn.iocoder.yudao.module.farm.service.devicealarm;
import cn.hutool.core.util.IdUtil;
import cn.iocoder.yudao.framework.common.pojo.PageResult;
import cn.iocoder.yudao.framework.security.core.LoginUser;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import cn.iocoder.yudao.module.farm.dal.dataobject.devicealarm.DeviceAlarmDO;
import cn.iocoder.yudao.module.farm.dal.dataobject.deviceinfo.DeviceInfoDO;
import cn.iocoder.yudao.module.farm.dal.mysql.devicealarm.DeviceAlarmMapper;
import cn.iocoder.yudao.module.farm.dal.mysql.deviceinfo.DeviceInfoMapper;
import cn.iocoder.yudao.module.farm.enums.DeviceAlarmTypeEnum;
import jakarta.annotation.Resource;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.validation.annotation.Validated;
import java.time.LocalDateTime;
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.DEVICE_ALARM_NOT_EXISTS;
/**
* 设备报警 Service 实现类
*
* @author 芋道源码
*/
@Service
@Validated
@Slf4j
public class DeviceAlarmServiceImpl implements DeviceAlarmService {
@Resource
private DeviceAlarmMapper deviceAlarmMapper;
@Resource
private DeviceInfoMapper deviceInfoMapper;
@Override
@Transactional(rollbackFor = Exception.class)
public String createDeviceAlarm(DeviceAlarmDO deviceAlarm) {
LoginUser user = getLoginUser();
// 设置基本信息
deviceAlarm.setId(IdUtil.fastUUID());
deviceAlarm.setTenantId(String.valueOf(user.getTenantId()));
deviceAlarm.setStatus(0); // 未处理
deviceAlarm.setPushStatus(0); // 未推送
deviceAlarm.setPushCount(0);
// 设置报警类型名称
DeviceAlarmTypeEnum alarmTypeEnum = DeviceAlarmTypeEnum.getByType(deviceAlarm.getAlarmType());
if (alarmTypeEnum != null) {
deviceAlarm.setAlarmTypeName(alarmTypeEnum.getName());
}
// 插入数据库
deviceAlarmMapper.insert(deviceAlarm);
log.info("[createDeviceAlarm][创建设备报警成功][deviceId: {}][alarmType: {}]",
deviceAlarm.getDeviceId(), deviceAlarm.getAlarmType());
return deviceAlarm.getId();
}
@Override
public void updateDeviceAlarm(DeviceAlarmDO deviceAlarm) {
// 校验存在
validateDeviceAlarmExists(deviceAlarm.getId());
// 更新
deviceAlarmMapper.updateById(deviceAlarm);
}
@Override
public void deleteDeviceAlarm(String id) {
// 校验存在
validateDeviceAlarmExists(id);
// 删除
deviceAlarmMapper.deleteById(id);
}
@Override
public void deleteDeviceAlarmListByIds(List<String> ids) {
deviceAlarmMapper.deleteBatchIds(ids);
}
@Override
public DeviceAlarmDO getDeviceAlarm(String id) {
return deviceAlarmMapper.selectById(id);
}
@Override
public PageResult<DeviceAlarmDO> getDeviceAlarmPage(Object pageReqVO) {
// TODO: 实现分页查询逻辑
return new PageResult<>();
}
@Override
public List<DeviceAlarmDO> getUnprocessedAlarmsByDeviceId(String deviceId) {
LambdaQueryWrapper<DeviceAlarmDO> wrapper = new LambdaQueryWrapper<DeviceAlarmDO>()
.eq(DeviceAlarmDO::getDeviceId, deviceId)
.eq(DeviceAlarmDO::getStatus, 0)
.orderByDesc(DeviceAlarmDO::getCreateTime);
return deviceAlarmMapper.selectList(wrapper);
}
@Override
public List<DeviceAlarmDO> getUnprocessedAlarmsByBaseId(String baseId) {
LambdaQueryWrapper<DeviceAlarmDO> wrapper = new LambdaQueryWrapper<DeviceAlarmDO>()
.eq(DeviceAlarmDO::getBaseId, baseId)
.eq(DeviceAlarmDO::getStatus, 0)
.orderByDesc(DeviceAlarmDO::getCreateTime);
return deviceAlarmMapper.selectList(wrapper);
}
@Override
@Transactional(rollbackFor = Exception.class)
public void processDeviceAlarm(String id, String processUserId, String processUserName, String processRemark) {
DeviceAlarmDO deviceAlarm = validateDeviceAlarmExists(id);
deviceAlarm.setStatus(1); // 已处理
deviceAlarm.setProcessTime(LocalDateTime.now());
deviceAlarm.setProcessUserId(processUserId);
deviceAlarm.setProcessUserName(processUserName);
deviceAlarm.setProcessRemark(processRemark);
deviceAlarmMapper.updateById(deviceAlarm);
log.info("[processDeviceAlarm][处理设备报警成功][id: {}][processUser: {}]", id, processUserName);
}
@Override
@Transactional(rollbackFor = Exception.class)
public void ignoreDeviceAlarm(String id, String processUserId, String processUserName, String processRemark) {
DeviceAlarmDO deviceAlarm = validateDeviceAlarmExists(id);
deviceAlarm.setStatus(2); // 已忽略
deviceAlarm.setProcessTime(LocalDateTime.now());
deviceAlarm.setProcessUserId(processUserId);
deviceAlarm.setProcessUserName(processUserName);
deviceAlarm.setProcessRemark(processRemark);
deviceAlarmMapper.updateById(deviceAlarm);
log.info("[ignoreDeviceAlarm][忽略设备报警成功][id: {}][processUser: {}]", id, processUserName);
}
@Override
@Transactional(rollbackFor = Exception.class)
public void checkAndCreateAlarm(String deviceId, Integer alarmType, String alarmData) {
// 获取设备信息
DeviceInfoDO device = deviceInfoMapper.selectById(deviceId);
if (device == null) {
log.warn("[checkAndCreateAlarm][设备不存在][deviceId: {}]", deviceId);
return;
}
// 检查是否已有相同类型的未处理报警
LambdaQueryWrapper<DeviceAlarmDO> existsWrapper = new LambdaQueryWrapper<DeviceAlarmDO>()
.eq(DeviceAlarmDO::getDeviceId, deviceId)
.eq(DeviceAlarmDO::getAlarmType, alarmType)
.eq(DeviceAlarmDO::getStatus, 0);
List<DeviceAlarmDO> existingAlarms = deviceAlarmMapper.selectList(existsWrapper);
if (!existingAlarms.isEmpty()) {
log.info("[checkAndCreateAlarm][已存在相同类型的未处理报警][deviceId: {}][alarmType: {}]", deviceId, alarmType);
return;
}
// 创建新报警
DeviceAlarmDO deviceAlarm = DeviceAlarmDO.builder()
.deviceId(deviceId)
.deviceName(device.getDeviceName())
.deviceUniqueId(device.getDeviceUniqueId())
.alarmType(alarmType)
.alarmLevel(determineAlarmLevel(alarmType))
.title(generateAlarmTitle(alarmType, device.getDeviceName()))
.content(generateAlarmContent(alarmType, device.getDeviceName()))
.alarmData(alarmData)
.baseId(device.getBaseId())
.baseName(device.getBaseName())
.plotId(device.getPlotId())
.plotName(device.getPlotName())
.build();
createDeviceAlarm(deviceAlarm);
}
/**
* 根据报警类型确定报警级别
*/
private Integer determineAlarmLevel(Integer alarmType) {
switch (alarmType) {
case 1:
case 2:
case 3:
return 2; // 中
case 4:
case 5:
case 6:
case 7:
return 3; // 高
case 8:
case 9:
case 10:
case 11:
return 2; // 中
case 12:
case 13:
return 4; // 紧急
case 14:
return 3; // 高
case 15:
return 1; // 低
default:
return 2; // 默认中
}
}
/**
* 生成报警标题
*/
private String generateAlarmTitle(Integer alarmType, String deviceName) {
DeviceAlarmTypeEnum alarmTypeEnum = DeviceAlarmTypeEnum.getByType(alarmType);
if (alarmTypeEnum != null) {
return deviceName + " - " + alarmTypeEnum.getName();
}
return deviceName + " - 设备报警";
}
/**
* 生成报警内容
*/
private String generateAlarmContent(Integer alarmType, String deviceName) {
DeviceAlarmTypeEnum alarmTypeEnum = DeviceAlarmTypeEnum.getByType(alarmType);
if (alarmTypeEnum != null) {
return "设备 " + deviceName + " 发生 " + alarmTypeEnum.getName() + ",请及时处理。";
}
return "设备 " + deviceName + " 发生报警,请及时处理。";
}
private DeviceAlarmDO validateDeviceAlarmExists(String id) {
DeviceAlarmDO deviceAlarm = deviceAlarmMapper.selectById(id);
if (deviceAlarm == null) {
throw exception(DEVICE_ALARM_NOT_EXISTS);
}
return deviceAlarm;
}
}

View File

@@ -1,27 +1,20 @@
package cn.iocoder.yudao.module.farm.service.messagecenter;
import cn.hutool.core.collection.CollUtil;
import cn.hutool.core.util.IdUtil;
import cn.iocoder.yudao.module.farm.controller.admin.messagecenter.vo.MessageCenterPageReqVO;
import cn.iocoder.yudao.module.farm.controller.admin.messagecenter.vo.MessageCenterSaveReqVO;
import org.springframework.stereotype.Service;
import jakarta.annotation.Resource;
import org.springframework.validation.annotation.Validated;
import org.springframework.transaction.annotation.Transactional;
import cn.iocoder.yudao.module.farm.dal.dataobject.messagecenter.MessageCenterDO;
import cn.iocoder.yudao.framework.common.pojo.PageResult;
import cn.iocoder.yudao.framework.common.pojo.PageParam;
import cn.iocoder.yudao.framework.common.util.object.BeanUtils;
import cn.iocoder.yudao.framework.security.core.LoginUser;
import cn.iocoder.yudao.module.farm.controller.admin.messagecenter.vo.MessageCenterPageReqVO;
import cn.iocoder.yudao.module.farm.controller.admin.messagecenter.vo.MessageCenterSaveReqVO;
import cn.iocoder.yudao.module.farm.dal.dataobject.messagecenter.MessageCenterDO;
import cn.iocoder.yudao.module.farm.dal.mysql.messagecenter.MessageCenterMapper;
import jakarta.annotation.Resource;
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.common.util.collection.CollectionUtils.convertList;
import static cn.iocoder.yudao.framework.common.util.collection.CollectionUtils.diffList;
import static cn.iocoder.yudao.framework.security.core.util.SecurityFrameworkUtils.getLoginUser;
import static cn.iocoder.yudao.module.farm.enums.ErrorCodeConstants.MESSAGE_CENTER_NOT_EXISTS;

View File

@@ -0,0 +1,98 @@
package cn.iocoder.yudao.module.farm.service.messagepush;
import cn.iocoder.yudao.module.farm.controller.admin.websocket.vo.FarmWebSocketMessageVO;
import cn.iocoder.yudao.module.farm.dal.dataobject.messagecenter.MessageCenterDO;
import java.util.List;
/**
* 农场消息推送 Service 接口
*
* @author 芋道源码
*/
public interface FarmMessagePushService {
/**
* 推送报警消息
*
* @param message 消息中心消息
*/
void pushAlarmMessage(MessageCenterDO message);
/**
* 推送预警消息
*
* @param message 消息中心消息
*/
void pushWarningMessage(MessageCenterDO message);
/**
* 推送普通消息
*
* @param message 消息中心消息
*/
void pushNormalMessage(MessageCenterDO message);
/**
* 推送消息到指定基地
*
* @param message 消息
* @param baseId 基地ID
*/
void pushMessageToBase(MessageCenterDO message, String baseId);
/**
* 推送消息到指定地块
*
* @param message 消息
* @param plotId 地块ID
*/
void pushMessageToPlot(MessageCenterDO message, String plotId);
/**
* 推送消息到指定用户
*
* @param message 消息
* @param userId 用户ID
*/
void pushMessageToUser(MessageCenterDO message, String userId);
/**
* 广播消息到所有在线用户
*
* @param message 消息
*/
void broadcastMessage(MessageCenterDO message);
/**
* 推送设备报警消息
*
* @param deviceId 设备ID
* @param deviceName 设备名称
* @param alarmType 报警类型
* @param alarmContent 报警内容
* @param baseId 基地ID
* @param plotId 地块ID
*/
void pushDeviceAlarm(String deviceId, String deviceName, Integer alarmType,
String alarmContent, String baseId, String plotId);
/**
* 推送天气预警消息
*
* @param warningContent 预警内容
* @param baseId 基地ID
* @param plotId 地块ID
*/
void pushWeatherWarning(String warningContent, String baseId, String plotId);
/**
* 推送系统通知消息
*
* @param title 标题
* @param content 内容
* @param baseId 基地ID
* @param plotId 地块ID
*/
void pushSystemNotification(String title, String content, String baseId, String plotId);
}

View File

@@ -0,0 +1,273 @@
package cn.iocoder.yudao.module.farm.service.messagepush;
import cn.hutool.core.util.IdUtil;
import cn.iocoder.yudao.framework.common.util.json.JsonUtils;
import cn.iocoder.yudao.framework.common.util.object.BeanUtils;
import cn.iocoder.yudao.framework.websocket.core.sender.WebSocketMessageSender;
import cn.iocoder.yudao.module.farm.controller.admin.messagecenter.vo.MessageCenterSaveReqVO;
import cn.iocoder.yudao.module.farm.controller.admin.websocket.vo.FarmWebSocketMessageVO;
import cn.iocoder.yudao.module.farm.dal.dataobject.messagecenter.MessageCenterDO;
import cn.iocoder.yudao.module.farm.enums.AlarmStatusEnum;
import cn.iocoder.yudao.module.farm.enums.MessageTypeEnum;
import cn.iocoder.yudao.module.farm.service.messagecenter.MessageCenterService;
import jakarta.annotation.Resource;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.time.LocalDateTime;
/**
* 农场消息推送 Service 实现类
*
* @author 芋道源码
*/
@Service
@Slf4j
public class FarmMessagePushServiceImpl implements FarmMessagePushService {
@Resource
private WebSocketMessageSender webSocketMessageSender;
@Resource
private MessageCenterService messageCenterService;
@Override
public void pushAlarmMessage(MessageCenterDO message) {
log.info("[pushAlarmMessage][推送报警消息][messageId: {}][title: {}]",
message.getId(), message.getTitle());
// 转换为WebSocket消息格式
FarmWebSocketMessageVO wsMessage = convertToWebSocketMessage(message);
// 推送消息
pushMessage(wsMessage);
}
@Override
public void pushWarningMessage(MessageCenterDO message) {
log.info("[pushWarningMessage][推送预警消息][messageId: {}][title: {}]",
message.getId(), message.getTitle());
// 转换为WebSocket消息格式
FarmWebSocketMessageVO wsMessage = convertToWebSocketMessage(message);
// 推送消息
pushMessage(wsMessage);
}
@Override
public void pushNormalMessage(MessageCenterDO message) {
log.info("[pushNormalMessage][推送普通消息][messageId: {}][title: {}]",
message.getId(), message.getTitle());
// 转换为WebSocket消息格式
FarmWebSocketMessageVO wsMessage = convertToWebSocketMessage(message);
// 推送消息
pushMessage(wsMessage);
}
@Override
public void pushMessageToBase(MessageCenterDO message, String baseId) {
log.info("[pushMessageToBase][推送消息到基地][baseId: {}][title: {}]",
baseId, message.getTitle());
// 转换为WebSocket消息格式
FarmWebSocketMessageVO wsMessage = convertToWebSocketMessage(message);
// 推送到指定基地的用户
// TODO: 实现按基地推送的逻辑
pushMessage(wsMessage);
}
@Override
public void pushMessageToPlot(MessageCenterDO message, String plotId) {
log.info("[pushMessageToPlot][推送消息到地块][plotId: {}][title: {}]",
plotId, message.getTitle());
// 转换为WebSocket消息格式
FarmWebSocketMessageVO wsMessage = convertToWebSocketMessage(message);
// 推送到指定地块的用户
// TODO: 实现按地块推送的逻辑
pushMessage(wsMessage);
}
@Override
public void pushMessageToUser(MessageCenterDO message, String userId) {
log.info("[pushMessageToUser][推送消息到用户][userId: {}][title: {}]",
userId, message.getTitle());
// 转换为WebSocket消息格式
FarmWebSocketMessageVO wsMessage = convertToWebSocketMessage(message);
// 推送到指定用户
webSocketMessageSender.send(1, Long.valueOf(userId), "farm", JsonUtils.toJsonString(wsMessage));
}
@Override
public void broadcastMessage(MessageCenterDO message) {
log.info("[broadcastMessage][广播消息][title: {}]", message.getTitle());
// 转换为WebSocket消息格式
FarmWebSocketMessageVO wsMessage = convertToWebSocketMessage(message);
// 广播消息
webSocketMessageSender.send(1, "farm", JsonUtils.toJsonString(wsMessage));
}
@Override
public void pushDeviceAlarm(String deviceId, String deviceName, Integer alarmType,
String alarmContent, String baseId, String plotId) {
log.info("[pushDeviceAlarm][推送设备报警][deviceId: {}][alarmType: {}]", deviceId, alarmType);
// 直接推送WebSocket消息,不创建消息中心记录
// 消息中心记录由调用方负责创建
FarmWebSocketMessageVO wsMessage = FarmWebSocketMessageVO.builder()
.messageType(MessageTypeEnum.ALARM.getType())
.messageTypeName(MessageTypeEnum.ALARM.getName())
.title(deviceName + " - 设备报警")
.content(alarmContent)
.messageId(IdUtil.fastUUID())
.deviceId(deviceId)
.deviceName(deviceName)
.baseId(baseId)
.plotId(plotId)
.alarmLevel(determineAlarmLevel(alarmType))
.messageTime(LocalDateTime.now())
.isRead(false)
.pushCount(1)
.build();
// 推送消息
pushMessage(wsMessage);
}
@Override
public void pushWeatherWarning(String warningContent, String baseId, String plotId) {
log.info("[pushWeatherWarning][推送天气预警][baseId: {}][plotId: {}]", baseId, plotId);
// 创建消息中心记录
MessageCenterDO message = createMessageCenterRecord(
MessageTypeEnum.WARNING.getType(),
"天气预警",
warningContent,
baseId,
plotId,
null,
null
);
// 推送预警消息
pushWarningMessage(message);
}
@Override
public void pushSystemNotification(String title, String content, String baseId, String plotId) {
log.info("[pushSystemNotification][推送系统通知][title: {}]", title);
// 创建消息中心记录
MessageCenterDO message = createMessageCenterRecord(
MessageTypeEnum.MESSAGE.getType(),
title,
content,
baseId,
plotId,
null,
null
);
// 推送普通消息
pushNormalMessage(message);
}
/**
* 转换为WebSocket消息格式
*/
private FarmWebSocketMessageVO convertToWebSocketMessage(MessageCenterDO message) {
FarmWebSocketMessageVO wsMessage = new FarmWebSocketMessageVO();
wsMessage.setMessageType(message.getMessageType());
wsMessage.setMessageTypeName(MessageTypeEnum.getByType(message.getMessageType()).getName());
wsMessage.setTitle(message.getTitle());
wsMessage.setContent(message.getContent());
wsMessage.setMessageId(message.getId());
wsMessage.setBaseId(message.getBaseId());
wsMessage.setBaseName(message.getBaseName());
wsMessage.setPlotId(message.getPlotId());
wsMessage.setPlotName(message.getPlotName());
wsMessage.setAlarmStatus(message.getAlarmStatus());
wsMessage.setMessageTime(message.getCreateTime());
wsMessage.setIsRead(message.getReadStatus() == 1);
wsMessage.setPushCount(message.getNumber());
return wsMessage;
}
/**
* 创建消息中心记录
*/
private MessageCenterDO createMessageCenterRecord(Integer messageType, String title, String content,
String baseId, String plotId, String deviceId, String deviceName) {
MessageCenterDO message = MessageCenterDO.builder()
.messageType(messageType)
.title(title)
.content(content)
.readStatus(0) // 未读
.alarmStatus(AlarmStatusEnum.UNRESOLVED.getStatus()) // 未解除
.number(0) // 推送次数
.baseId(baseId)
.plotId(plotId)
.build();
// 保存到消息中心
String messageCenter = messageCenterService.createMessageCenter(BeanUtils.toBean(message, MessageCenterSaveReqVO.class));// TODO: 需要创建对应的VO
return message.setId(messageCenter);
}
/**
* 根据报警类型确定报警级别
*/
private Integer determineAlarmLevel(Integer alarmType) {
// 根据报警类型确定级别
switch (alarmType) {
case 1: // 设备离线
case 2: // 电池电量低
case 3: // 信号强度弱
return 2; // 中
case 4: // 温度过高
case 5: // 温度过低
case 6: // 湿度过高
case 7: // 湿度过低
return 3; // 高
case 8: // 土壤湿度低
case 9: // 土壤湿度过高
case 10: // pH值异常
case 11: // 养分不足
return 2; // 中
case 12: // 病虫害检测
case 13: // 病害检测
return 4; // 紧急
case 14: // 天气预警
return 3; // 高
case 15: // 设备维护到期
return 1; // 低
default:
return 2; // 默认中
}
}
/**
* 推送消息
*/
private void pushMessage(FarmWebSocketMessageVO wsMessage) {
try {
// 广播消息到所有在线用户
webSocketMessageSender.send(1, "farm", JsonUtils.toJsonString(wsMessage));
log.info("[pushMessage][消息推送成功][messageId: {}][title: {}]",
wsMessage.getMessageId(), wsMessage.getTitle());
} catch (Exception e) {
log.error("[pushMessage][消息推送失败][messageId: {}][title: {}]",
wsMessage.getMessageId(), wsMessage.getTitle(), e);
}
}
}

View File

@@ -0,0 +1,74 @@
package cn.iocoder.yudao.module.farm.websocket;
import cn.iocoder.yudao.framework.websocket.core.listener.WebSocketMessageListener;
import cn.iocoder.yudao.framework.websocket.core.message.JsonWebSocketMessage;
import cn.iocoder.yudao.module.farm.controller.admin.websocket.vo.FarmWebSocketMessageVO;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import org.springframework.web.socket.WebSocketSession;
/**
* 农场WebSocket消息监听器
*
* @author 芋道源码
*/
@Component
@Slf4j
public class FarmWebSocketMessageListener implements WebSocketMessageListener<FarmWebSocketMessageVO> {
@Override
public String getType() {
return "farm";
}
@Override
public void onMessage(WebSocketSession session, FarmWebSocketMessageVO message) {
log.info("[onMessage][收到农场WebSocket消息][session: {}][message: {}]", session.getId(), message);
// 处理不同类型的消息
switch (message.getMessageType()) {
case 1: // 报警消息
handleAlarmMessage(session, message);
break;
case 2: // 预警消息
handleWarningMessage(session, message);
break;
case 3: // 普通消息
handleNormalMessage(session, message);
break;
default:
log.warn("[onMessage][未知的消息类型][messageType: {}]", message.getMessageType());
}
}
/**
* 处理报警消息
*/
private void handleAlarmMessage(WebSocketSession session, FarmWebSocketMessageVO message) {
log.info("[handleAlarmMessage][处理报警消息][messageId: {}][title: {}]",
message.getMessageId(), message.getTitle());
// 这里可以添加报警消息的特殊处理逻辑
// 比如记录日志、发送通知等
}
/**
* 处理预警消息
*/
private void handleWarningMessage(WebSocketSession session, FarmWebSocketMessageVO message) {
log.info("[handleWarningMessage][处理预警消息][messageId: {}][title: {}]",
message.getMessageId(), message.getTitle());
// 这里可以添加预警消息的特殊处理逻辑
}
/**
* 处理普通消息
*/
private void handleNormalMessage(WebSocketSession session, FarmWebSocketMessageVO message) {
log.info("[handleNormalMessage][处理普通消息][messageId: {}][title: {}]",
message.getMessageId(), message.getTitle());
// 这里可以添加普通消息的特殊处理逻辑
}
}