refactor(mqtt): 重构MQTT消息处理器并优化设备通信逻辑

- 添加AskDeviceDataFeignClient依赖注入以支持设备数据询问功能
- 移除废弃的/devTopic/{edgeId}主题订阅方法及相关注释代码
- 移除废弃的devOperation方法和devAccessOperation方法的注释实现
- 更新/dev/Data/{version}/{edgeId}主题处理逻辑,重新定义为文件传输功能
- 新增对/csUpgradeLogsFeignClient和askDeviceDataFeignClient的依赖注入
- 优化文件传输逻辑,添加升级进度通知和设备重启功能
- 移除旧的文件传输处理方法注释代码
- 添加新的/dev/Error/{edgeId}主题处理器用于记录装置异常事件
- 移除旧的异常事件处理方法注释代码
- 优化升级文件CRC校验和设备升级状态更新逻辑
This commit is contained in:
xy
2026-07-23 15:39:19 +08:00
parent 73536f1c28
commit 9e35193e10

View File

@@ -12,6 +12,7 @@ import com.github.tocrhz.mqtt.annotation.MqttSubscribe;
import com.github.tocrhz.mqtt.annotation.NamedValue; import com.github.tocrhz.mqtt.annotation.NamedValue;
import com.github.tocrhz.mqtt.annotation.Payload; import com.github.tocrhz.mqtt.annotation.Payload;
import com.github.tocrhz.mqtt.publisher.MqttPublisher; import com.github.tocrhz.mqtt.publisher.MqttPublisher;
import com.njcn.access.api.AskDeviceDataFeignClient;
import com.njcn.access.enums.AccessEnum; import com.njcn.access.enums.AccessEnum;
import com.njcn.access.enums.AccessResponseEnum; import com.njcn.access.enums.AccessResponseEnum;
import com.njcn.access.enums.TypeEnum; import com.njcn.access.enums.TypeEnum;
@@ -102,10 +103,18 @@ public class MqttMessageHandler {
private final CsDeviceRegistryFeignClient csDeviceRegistryFeignClient; private final CsDeviceRegistryFeignClient csDeviceRegistryFeignClient;
private final ExecutorService mqttMessageExecutor; private final ExecutorService mqttMessageExecutor;
private final CsEdDataFeignClient csEdDataFeignClient; private final CsEdDataFeignClient csEdDataFeignClient;
private static Integer mid = 1; private final CsUpgradeLogsFeignClient csUpgradeLogsFeignClient;
private final AskDeviceDataFeignClient askDeviceDataFeignClient;
@Autowired @Autowired
Validator validator; Validator validator;
/**
* 获取主题
* @param topic
* @param message
* @param nDid
* @param payload
*/
@MqttSubscribe(value = "/Dev/DevTopic/{edgeId}",qos = 1) @MqttSubscribe(value = "/Dev/DevTopic/{edgeId}",qos = 1)
public void devTopic(String topic, MqttMessage message, @NamedValue("edgeId") String nDid, @Payload String payload){ public void devTopic(String topic, MqttMessage message, @NamedValue("edgeId") String nDid, @Payload String payload){
Gson gson = new Gson(); Gson gson = new Gson();
@@ -158,56 +167,14 @@ public class MqttMessageHandler {
} }
} }
// @MqttSubscribe(value = "/Dev/DevTopic/{edgeId}",qos = 1) /**
// public void devTopic(String topic, MqttMessage message, @NamedValue("edgeId") String nDid, @Payload String payload){ * 装置注册应答
// //业务流程开始 * 1.收到注册信息,修改装置出厂表,装置的状态,调整为注册;然后开始接入流程
// Gson gson = new Gson(); * 2.询问当前装置类型的模板。有则完成接入;没有则告警出来,需要人工手动上传模板信息
// ReqAndResParam.Res res = gson.fromJson(new String(message.getPayload(), StandardCharsets.UTF_8), ReqAndResParam.Res.class); * @param topic
// //日志记录 * @param message
// LogMessage logDto = new LogMessage(); * @param payload
// logDto.setUserIndex("系统"); */
// logDto.setLoginName("系统");
// logDto.setOperate("系统端收到装置端"+nDid+"发送的主题信息code = " + res.getCode());
// logDto.setResult(1);
// logMessageTemplate.sendMember(logDto);
// //检验传递的参数是否准确
// Set<ConstraintViolation<ReqAndResParam.Res>> validate = validator.validate(res);
//// validate.forEach(constraintViolation -> {
//// System.out.println(constraintViolation.getMessage());
//// });
// if (Objects.equals(res.getCode(),AccessEnum.SUCCESS.getCode())){
// if (Objects.equals(res.getType(), Integer.parseInt(TypeEnum.TYPE_16.getCode()))){
// List<CsTopic> list = new ArrayList<>();
// Map<String,List<String>> map = (Map<String,List<String>>)res.getMsg();
// List<String> topicList = map.get("Topic");
// topicList.forEach(item->{
// CsTopic csTopic = new CsTopic();
// csTopic.setNDid(nDid);
// csTopic.setTopic(item);
// String version = item.split("/")[3];
// if (version.startsWith("V")){
// csTopic.setVersion(version);
// }
// list.add(csTopic);
// });
// csTopicService.addTopic(nDid,list);
// String version = list.stream().map(CsTopic::getVersion).filter(Objects::nonNull).findFirst().orElse("V1");
// redisUtil.saveByKeyWithExpire(nDid +":version",version,30L);
// logMessageTemplate.sendMember(logDto);
// } else {
// logDto.setResult(0);
// logDto.setFailReason(AccessResponseEnum.MESSAGE_TYPE_ERROR.getMessage());
// logMessageTemplate.sendMember(logDto);
// //log.info(AccessResponseEnum.MESSAGE_TYPE_ERROR.getMessage());
// }
// } else {
// logDto.setResult(0);
// logDto.setFailReason(AccessResponseEnum.RESPONSE_ERROR.getMessage());
// logMessageTemplate.sendMember(logDto);
// //log.info(AccessResponseEnum.RESPONSE_ERROR.getMessage());
// }
// }
@MqttSubscribe(value = "/Dev/DevReg/{edgeId}",qos = 1) @MqttSubscribe(value = "/Dev/DevReg/{edgeId}",qos = 1)
public void devOperation(String topic, MqttMessage message, @NamedValue("edgeId") String nDid, @Payload String payload){ public void devOperation(String topic, MqttMessage message, @NamedValue("edgeId") String nDid, @Payload String payload){
log.info("收到注册应答响应--->{}", nDid); log.info("收到注册应答响应--->{}", nDid);
@@ -256,58 +223,14 @@ public class MqttMessageHandler {
} }
} }
// /** /**
// * 装置注册应答 * 设备响应
// * 1.收到注册信息,修改装置出厂表,装置的状态,调整为注册;然后开始接入流程 * @param topic
// * 2.询问当前装置类型的模板。有则完成接入;没有则告警出来,需要人工手动上传模板信息 * @param message
// * @param topic * @param version
// * @param message * @param nDid
// * @param payload * @param payload
// */ */
// @MqttSubscribe(value = "/Dev/DevReg/{edgeId}",qos = 1)
// public void devOperation(String topic, MqttMessage message, @NamedValue("edgeId") String nDid, @Payload String payload){
// log.info("收到注册应答响应--->{}", nDid);
// //业务处理
// Gson gson = new Gson();
// ReqAndResDto.Res res = gson.fromJson(new String(message.getPayload(), StandardCharsets.UTF_8), ReqAndResDto.Res.class);
// //日志记录
// LogMessage logDto = new LogMessage();
// logDto.setUserIndex("系统");
// logDto.setLoginName("系统");
// logDto.setOperate("系统端收到装置端"+nDid+"注册应答响应code = " + res.getCode());
// logDto.setResult(1);
// logMessageTemplate.sendMember(logDto);
// if (Objects.equals(res.getCode(),AccessEnum.SUCCESS.getCode())){
// if (Objects.equals(res.getType(),Integer.parseInt(TypeEnum.TYPE_17.getCode()))){
// //询问模板数据
// ReqAndResDto.Req reqAndResParam = new ReqAndResDto.Req();
// reqAndResParam.setMid(1);
// reqAndResParam.setDid(0);
// reqAndResParam.setPri(AccessEnum.FIRST_CHANNEL.getCode());
// reqAndResParam.setType(Integer.parseInt(TypeEnum.TYPE_3.getCode()));
// reqAndResParam.setExpire(-1);
// String version = csTopicService.getVersion(nDid);
// publisher.send("/Pfm/DevCmd/"+version+"/"+nDid,new Gson().toJson(reqAndResParam),1,false);
// //记录日志
// logDto.setUserIndex("系统");
// logDto.setLoginName("系统");
// logDto.setOperate("注册阶段:系统端向装置端"+nDid+"发送询问模板请求");
// logDto.setResult(1);
// logMessageTemplate.sendMember(logDto);
// } else {
// logDto.setResult(0);
// logDto.setFailReason(AccessResponseEnum.MESSAGE_TYPE_ERROR.getMessage());
// logMessageTemplate.sendMember(logDto);
// //log.info(AccessResponseEnum.MESSAGE_TYPE_ERROR.getMessage());
// }
// } else {
// logDto.setResult(0);
// logDto.setFailReason(AccessResponseEnum.REGISTER_RESPONSE_ERROR.getMessage());
// logMessageTemplate.sendMember(logDto);
// //log.info(AccessResponseEnum.REGISTER_RESPONSE_ERROR.getMessage());
// }
// }
@MqttSubscribe(value = "/Pfm/DevRsp/{version}/{edgeId}",qos = 1) @MqttSubscribe(value = "/Pfm/DevRsp/{version}/{edgeId}",qos = 1)
public void devAccessOperation(String topic, MqttMessage message, @NamedValue("version") String version, @NamedValue("edgeId") String nDid, @Payload String payload) { public void devAccessOperation(String topic, MqttMessage message, @NamedValue("version") String version, @NamedValue("edgeId") String nDid, @Payload String payload) {
LogMessage logDto = new LogMessage(); LogMessage logDto = new LogMessage();
@@ -564,302 +487,6 @@ public class MqttMessageHandler {
} }
} }
// /**
// * 设备响应
// * @param topic
// * @param message
// * @param version
// * @param nDid
// * @param payload
// */
// @MqttSubscribe(value = "/Pfm/DevRsp/{version}/{edgeId}",qos = 1)
// public void devAccessOperation(String topic, MqttMessage message, @NamedValue("version") String version, @NamedValue("edgeId") String nDid, @Payload String payload) throws InterruptedException {
// //日志实体
// LogMessage logDto = new LogMessage();
// logDto.setUserIndex("系统");
// logDto.setLoginName("系统");
// logDto.setResult(1);
// //业务处理
// Gson gson = new Gson();
// ReqAndResDto.Res res = gson.fromJson(new String(message.getPayload(), StandardCharsets.UTF_8), ReqAndResDto.Res.class);
// //redisUtil.saveByKeyWithExpire("devResponse",res.getCode(),5L);
// if (Objects.equals(res.getCode(),AccessEnum.SUCCESS.getCode())) {
// switch (res.getType()){
// /**
// * 装置类型模板应答
// * 1.判断网关的类型
// * 2.直联设备的DevCfg和DevMod是以直联设备为准上送平台端平台端保存。通过校验DevMod模板信息来从平台端模板池中选取对应的模板如果找不到匹配模板需告警提示人工干预处理。
// * 3.平台端需读取装置的DevMod来判断网关支持的设备模板包含设备型号和模板版本根据app提交的接入子设备DID匹配数据模板型号及版本生成DevCfg下发给网关网关根据下发信息生成就地设备点表。
// */
// case 4611:
// logDto.setOperate("系统端收到装置端"+nDid+"模板应答code = " + res.getCode());
// logMessageTemplate.sendMember(logDto);
// //log.info("{},装置模板应答,应答code {}",nDid,res.getCode());
// ModelDto modelDto = gson.fromJson(new String(message.getPayload(), StandardCharsets.UTF_8), ModelDto.class);
// List<DevModInfoDto> list = modelDto.getMsg().getDevMod();
// List<DevCfgDto> list2 = modelDto.getMsg().getDevCfg();
// if (CollectionUtils.isEmpty(list)) {
// //log.error(AccessResponseEnum.MODEL_VERSION_ERROR.getMessage());
// logDto.setOperate("查看装置端模板报文数据");
// logDto.setResult(0);
// logDto.setFailReason(AccessResponseEnum.MODEL_VERSION_ERROR.getMessage());
// logMessageTemplate.sendMember(logDto);
// //有异常删除缓存的模板信息
// redisUtil.delete(AppRedisKey.MODEL + nDid);
// log.error("{}{}", nDid, AccessResponseEnum.MODEL_VERSION_ERROR.getMessage());
// throw new BusinessException(AccessResponseEnum.MODEL_VERSION_ERROR);
// }
// //校验前置传递的装置模板库中是否存在
// List<CsModelDto> modelList = new ArrayList<>();
// list.forEach(item->{
// Integer did = null;
// for (DevCfgDto item2 : list2) {
// if (Objects.equals(item.getDevType(),item2.getDevType())){
// did = item2.getDid();
// }
// }
// CsModelDto csModelDto = new CsModelDto();
// CsDevModelPO po = devModelFeignClient.findModel(item.getDevType(),item.getVersionNo(),item.getVersionDate()).getData();
// if (Objects.isNull(po)){
// //log.error(AccessResponseEnum.MODEL_NO_FIND.getMessage());
// logDto.setOperate(nDid + "查询系统中是否存在模板");
// logDto.setResult(0);
// logDto.setFailReason(AccessResponseEnum.MODEL_NO_FIND.getMessage());
// logMessageTemplate.sendMember(logDto);
// //有异常删除缓存的模板信息
// redisUtil.delete(AppRedisKey.MODEL + nDid);
// throw new BusinessException(AccessResponseEnum.MODEL_NO_FIND);
// }
// if (Objects.equals(po.getType(),0)){
// List<CsDataSet> dataSetList = dataSetFeignClient.getModuleDataSet(po.getId()).getData();
// if (CollectionUtils.isEmpty(dataSetList)){
// logDto.setOperate("查看APF模块个数");
// logDto.setResult(0);
// logDto.setFailReason(AccessResponseEnum.MODULE_NUMBER_IS_NULL.getMessage());
// logMessageTemplate.sendMember(logDto);
// //有异常删除缓存的模板信息
// redisUtil.delete(AppRedisKey.MODEL + nDid);
// throw new BusinessException(AccessResponseEnum.MODULE_NUMBER_IS_NULL);
// }
// csModelDto.setModuleNumber(dataSetList.size());
// }
// csModelDto.setDevType(po.getDevTypeName());
// csModelDto.setModelId(po.getId());
// csModelDto.setDid(did);
// csModelDto.setType(po.getType());
// modelList.add(csModelDto);
// });
// //存储模板id
// String key2 = AppRedisKey.MODEL + nDid;
// redisUtil.saveByKey(key2,modelList);
// //存储监测点模板信息,用于界面回显
// List<String> modelId = modelList.stream().map(CsModelDto::getModelId).collect(Collectors.toList());
// List<CsLineModel> lineList = csLineModelService.getMonitorNumByModelId(modelId);
// String key = AppRedisKey.LINE + nDid;
// redisUtil.saveByKeyWithExpire(key,lineList,600L);
// logMessageTemplate.sendMember(logDto);
// break;
// case 4613:
// logDto.setOperate("系统端收到装置端"+nDid+"接入应答code = " + res.getCode());
// logDto.setResult(1);
// logMessageTemplate.sendMember(logDto);
// //log.info("{},收到接入应答响应,应答code {}",nDid,res.getCode());
// if (Objects.equals(res.getCode(),AccessEnum.SUCCESS.getCode())){
// int mid = 1;
// //修改装置状态
// csEquipmentDeliveryService.updateStatusBynDid(nDid,AccessEnum.ACCESS.getCode(),null,null);
// csEquipmentDeliveryService.updateRunStatusBynDid(nDid,AccessEnum.ONLINE.getCode());
// //记录设备上线
// PqsCommunicateDto dto = new PqsCommunicateDto();
// dto.setTime(LocalDateTime.now().format(DateTimeFormatter.ofPattern(DatePattern.NORM_DATETIME_PATTERN)));
// dto.setDevId(nDid);
// dto.setType(1);
// dto.setDescription("通讯正常");
// csCommunicateFeignClient.insertion(dto);
// //询问设备软件信息
// askDevData(nDid,version,1,mid);
// //更新治理监测点信息和设备容量
// askDevData(nDid,version,2,(res.getMid()+1));
// //更新电网侧、负载侧监测点信息
// askDevData(nDid,version,3,(res.getMid()+1));
// //接入后系统重置装置心跳
// heartbeatService.receiveHeartbeat(nDid);
// //修改redis的mid
// redisUtil.saveByKey(AppRedisKey.DEVICE_MID + nDid,1);
// //接入成功标识
// redisUtil.saveByKeyWithExpire("online" + nDid,"online",10L);
// //录波任务倒计时
// redisUtil.saveByKeyWithExpire("startFile:" + nDid,null,60L);
// } else {
// //log.info(AccessResponseEnum.ACCESS_RESPONSE_ERROR.getMessage());
// logDto.setResult(0);
// logDto.setFailReason(AccessResponseEnum.ACCESS_RESPONSE_ERROR.getMessage());
// logMessageTemplate.sendMember(logDto);
// throw new BusinessException(AccessResponseEnum.ACCESS_RESPONSE_ERROR);
// }
// logMessageTemplate.sendMember(logDto);
// break;
// case 4614:
// RspDataDto rspDataDto = JSON.parseObject(JSON.toJSONString(res.getMsg()), RspDataDto.class);
// if (!Objects.isNull(rspDataDto.getDataType())) {
// switch (rspDataDto.getDataType()){
// case 1:
// logDto.setOperate("系统端收到装置端"+nDid+"更新设备软件信息报文code = " + res.getCode());
// logDto.setResult(1);
// RspDataDto.SoftInfo softInfo = JSON.parseObject(JSON.toJSONString(rspDataDto.getDataArray()), RspDataDto.SoftInfo.class);
// //记录设备软件信息
// CsSoftInfoPO csSoftInfoPo = new CsSoftInfoPO();
// BeanUtils.copyProperties(softInfo,csSoftInfoPo);
// String id = IdUtil.fastSimpleUUID();
// csSoftInfoPo.setId(id);
// DateTimeFormatter formatter = new DateTimeFormatterBuilder()
// // 优先尝试紧凑格式
// .optionalStart().appendPattern("yyyyMMdd").optionalEnd()
// // 再尝试带横线格式
// .optionalStart().appendPattern("yyyy-MM-dd").optionalEnd()
// // 默认时间部分
// .parseDefaulting(ChronoField.HOUR_OF_DAY, 0)
// .parseDefaulting(ChronoField.MINUTE_OF_HOUR, 0)
// .parseDefaulting(ChronoField.SECOND_OF_MINUTE, 0)
// .toFormatter();
// LocalDateTime localDateTime = LocalDateTime.parse(softInfo.getAppDate(), formatter);
// assertThat(localDateTime).isNotNull();
// csSoftInfoPo.setAppDate(localDateTime);
// csSoftInfoFeignClient.saveSoftInfo(csSoftInfoPo);
// //更新设备软件id 先看是否存在软件信息,删除 然后在录入
// CsEquipmentDeliveryPO po = equipmentFeignClient.findDevByNDid(nDid).getData();
// String soft = po.getSoftinfoId();
// if (StringUtil.isNotBlank(soft)){
// csSoftInfoFeignClient.removeSoftInfo(soft);
// }
// equipmentFeignClient.updateSoftInfo(nDid,csSoftInfoPo.getId());
// logMessageTemplate.sendMember(logDto);
// break;
// case 2:
// List<RspDataDto.LdevInfo> devInfo = JSON.parseArray(JSON.toJSONString(rspDataDto.getDataArray()), RspDataDto.LdevInfo.class);
// if (CollectionUtil.isNotEmpty(devInfo)){
// if (Objects.equals(res.getDid(),1)){
// redisUtil.saveByKeyWithExpire("lineInfo:"+nDid,devInfo,30L);
// List<CsDevCapacityPO> list3 = new ArrayList<>();
// boolean hasZeroClDid = devInfo.stream().anyMatch(item -> item.getClDid() == 0);
// //治理设备
// if (hasZeroClDid) {
// logDto.setOperate("系统端收到装置端"+nDid+"更新APF容量报文code = " + res.getCode());
// logDto.setResult(1);
// devInfo.forEach(item->{
// if (Objects.equals(item.getClDid(),0)){
// updateLineInfo(nDid,item);
// }
// //2.录入各个模块设备容量
// CsDeviceRegistry csDeviceRegistry = csDeviceRegistryFeignClient.queryByCurrentNdidAndClDid(nDid, 0).getData();
// String lineId;
// if (Objects.isNull(csDeviceRegistry)) {
// List<CsLinePO> lines = csLineFeignClient.findByNdid(nDid).getData();
// if (CollectionUtil.isEmpty(lines)) {
// throw new BusinessException("通过装置NDID获取监测点信息失败");
// }
// CsLinePO line = lines.stream().filter(item2->Objects.equals(item2.getClDid(),0)).findFirst().orElse(null);
// if (Objects.isNull(line)) {
// throw new BusinessException("通过逻辑子设备ID获取监测点信息失败");
// }
// lineId = line.getLineId();
// } else {
// lineId = csDeviceRegistry.getId();
// }
// CsDevCapacityPO csDevCapacity = new CsDevCapacityPO();
// csDevCapacity.setLineId(lineId);
// csDevCapacity.setCldid(item.getClDid());
// csDevCapacity.setCapacity(Objects.isNull(item.getCapacityA())?0.0:item.getCapacityA());
// list3.add(csDevCapacity);
// });
// }
// //其余设备
// else {
// logDto.setOperate("系统端收到装置端"+nDid+"更新监测点台账报文code = " + res.getCode());
// logDto.setResult(1);
// devInfo.forEach(item->{
// updateLineInfo(nDid,item);
// });
// }
// if (CollectionUtil.isNotEmpty(list3)) {
// devCapacityFeignClient.addList(list3);
// //3.更新设备模块个数
// equipmentFeignClient.updateModuleNumber(nDid,(devInfo.size()-1));
// }
// } else if (Objects.equals(res.getDid(),2)) {
// logDto.setOperate("系统端收到装置端"+nDid+"更新电网侧、负载侧监测点信息报文code = " + res.getCode());
// logDto.setResult(1);
// //1.更新电网侧、负载侧监测点相关信息
// devInfo.forEach(item->{
// updateLineInfo(nDid,item);
// });
// }
// }
// logMessageTemplate.sendMember(logDto);
// break;
// case 15:
// logDto.setOperate("系统端收到装置端"+nDid+"更新设备软件信息报文code = " + res.getCode());
// logDto.setResult(1);
// JSONObject jsonObject = JSONObject.parseObject(JSON.toJSONString(res));
// AppAutoDataMessage appAutoDataMessage = JSONObject.toJavaObject(jsonObject, AppAutoDataMessage.class);
// appAutoDataMessage.setId(nDid);
// rtFeignClient.apfRtAnalysis(appAutoDataMessage);
// logMessageTemplate.sendMember(logDto);
// break;
// case 48:
// logDto.setOperate("系统端收到装置端"+nDid+"询问项目列表报文code = " + res.getCode());
// logDto.setResult(1);
// List<RspDataDto.ProjectInfo> projectInfoList = JSON.parseArray(JSON.toJSONString(rspDataDto.getDataArray()), RspDataDto.ProjectInfo.class);
// CsDeviceRegistry csDeviceRegistry = csDeviceRegistryFeignClient.queryByCurrentNdidAndClDid(nDid, rspDataDto.getClDid()).getData();
// String lineId;
// if (Objects.isNull(csDeviceRegistry)) {
// List<CsLinePO> lines = csLineFeignClient.findByNdid(nDid).getData();
// if (CollectionUtil.isEmpty(lines)) {
// throw new BusinessException("通过装置NDID获取监测点信息失败");
// }
// CsLinePO line = lines.stream().filter(item->Objects.equals(item.getClDid(),rspDataDto.getClDid())).findFirst().orElse(null);
// if (Objects.isNull(line)) {
// throw new BusinessException("通过逻辑子设备ID获取监测点信息失败");
// }
// lineId = line.getLineId();
// } else {
// lineId = csDeviceRegistry.getId();
// }
// String key3 = AppRedisKey.PROJECT_INFO + lineId;
// redisUtil.saveByKeyWithExpire(key3,projectInfoList,60L);
// logMessageTemplate.sendMember(logDto);
// break;
// default:
// break;
// }
// }
// break;
// case 4663:
// if (Objects.equals(res.getCode(),AccessEnum.SUCCESS.getCode())){
// String key4 = AppRedisKey.CONTROL + nDid;
// redisUtil.saveByKeyWithExpire(key4,"success",10L);
// }
// break;
// default:
// break;
// }
// } else if (Objects.equals(res.getCode(),AccessEnum.START_CHANNEL.getCode())) {
// logDto.setOperate(AccessEnum.START_CHANNEL.getMessage() + "系统等待5s");
// logDto.setResult(1);
// logMessageTemplate.sendMember(logDto);
// Thread.sleep(5000);
// } else {
// String result = getEnum(res.getCode());
// //log.info(result);
// logDto.setOperate("装置响应");
// logDto.setResult(0);
// logDto.setFailReason(result);
// logMessageTemplate.sendMember(logDto);
// throw new BusinessException(result);
// }
// }
public void updateLineInfo(String nDid,RspDataDto.LdevInfo item) { public void updateLineInfo(String nDid,RspDataDto.LdevInfo item) {
CsDeviceRegistry csDeviceRegistry = csDeviceRegistryFeignClient.queryByCurrentNdidAndClDid(nDid, item.getClDid()).getData(); CsDeviceRegistry csDeviceRegistry = csDeviceRegistryFeignClient.queryByCurrentNdidAndClDid(nDid, item.getClDid()).getData();
String lineId; String lineId;
@@ -891,6 +518,15 @@ public class MqttMessageHandler {
overLimitWlMapper.insert(overlimit); overLimitWlMapper.insert(overlimit);
} }
/**
* 装置心跳 && 主动数据上送
* fixme 这边由于接收文件数据时间跨度会很长,途中有其他请求进来会中断之前的程序,目前是记录中断的位置,等处理完成再继续请求接收文件
* @param topic
* @param message
* @param version
* @param nDid
* @param payload
*/
@MqttSubscribe(value = "/Dev/Data/{version}/{edgeId}",qos = 1) @MqttSubscribe(value = "/Dev/Data/{version}/{edgeId}",qos = 1)
public void devHeartBeat(String topic, MqttMessage message, @NamedValue("version") String version, @NamedValue("edgeId") String nDid, @Payload String payload) { public void devHeartBeat(String topic, MqttMessage message, @NamedValue("version") String version, @NamedValue("edgeId") String nDid, @Payload String payload) {
Gson gson = new Gson(); Gson gson = new Gson();
@@ -1001,102 +637,14 @@ public class MqttMessageHandler {
} }
} }
/** /**
* 装置心跳 && 主动数据上送 * 文件传输
* fixme 这边由于接收文件数据时间跨度会很长,途中有其他请求进来会中断之前的程序,目前是记录中断的位置,等处理完成再继续请求接收文件
* @param topic * @param topic
* @param message * @param message
* @param version * @param version
* @param nDid * @param nDid
* @param payload * @param payload
*/ */
// @MqttSubscribe(value = "/Dev/Data/{version}/{edgeId}",qos = 1)
// public void devHeartBeat(String topic, MqttMessage message, @NamedValue("version") String version, @NamedValue("edgeId") String nDid, @Payload String payload) {
// //解析数据
// Gson gson = new Gson();
// ReqAndResDto.Req res = gson.fromJson(new String(message.getPayload(), StandardCharsets.UTF_8), ReqAndResDto.Req.class);
// //响应请求
// switch (res.getType()){
// case 4865:
// heartbeatService.receiveHeartbeat(nDid);
// //有心跳,则将装置改成在线
// //csEquipmentDeliveryService.updateRunStatusBynDid(nDid,AccessEnum.ONLINE.getCode());
// //处理心跳 判断设备是否接入,如果设备已经接入则响应,不然忽略
// CsEquipmentDeliveryPO po = equipmentFeignClient.findDevByNDid(nDid).getData();
// if (Objects.nonNull(po)) {
// if (po.getUsageStatus() == 1 && po.getRunStatus() == 2 && po.getStatus() == 3) {
// ReqAndResDto.Res reqAndResParam = new ReqAndResDto.Res();
// reqAndResParam.setMid(res.getMid());
// reqAndResParam.setDid(0);
// reqAndResParam.setPri(AccessEnum.FIRST_CHANNEL.getCode());
// reqAndResParam.setType(Integer.parseInt(TypeEnum.TYPE_29.getCode()));
// reqAndResParam.setCode(200);
// //fixme 前置处理的时间应该是UTC时间所以需要加8小时。
// String json = "{Time:"+(System.currentTimeMillis()/1000+8*3600)+"}";
// net.sf.json.JSONObject jsonObject = net.sf.json.JSONObject.fromObject(json);
// reqAndResParam.setMsg(jsonObject);
// publisher.send("/Dev/DataRsp/"+version+"/"+nDid,gson.toJson(reqAndResParam),1,false);
// //处理业务逻辑
// Object object = res.getMsg();
// if (!Objects.isNull(object)){
// List<String> abnormalList = new ArrayList<>();
// if (object instanceof ArrayList<?>){
// abnormalList.addAll((List<String>) object);
// }
// //todo APF设备不存在逻辑设备掉线的情况网关下的设备会存在
// abnormalList.forEach(item->{
// System.out.println("异常设备ID"+item);
// });
// }
// }
// }
// break;
// case 4866:
// AutoDataDto dataDto = gson.fromJson(new String(message.getPayload(), StandardCharsets.UTF_8), AutoDataDto.class);
// //mid大于0则需要应答设备侧
// if (dataDto.getMid() > 0){
// ReqAndResDto.Res response = new ReqAndResDto.Res();
// response.setMid(dataDto.getMid());
// response.setDid(dataDto.getDid());
// response.setPri(AccessEnum.FIRST_CHANNEL.getCode());
// response.setType(Integer.parseInt(TypeEnum.TYPE_15.getCode()));
// response.setCode(200);
// publisher.send("/Dev/DataRsp/"+version+"/"+nDid,new Gson().toJson(response),1,false);
// }
// //判断事件类型
// switch (dataDto.getMsg().getDataAttr()) {
// //暂态事件、录波处理、工程信息
// case 0:
// EventDto eventDto = gson.fromJson(new String(message.getPayload(), StandardCharsets.UTF_8), EventDto.class);
// JSONObject jsonObject0 = JSONObject.parseObject(JSON.toJSONString(eventDto));
// AppEventMessage appEventMessage = JSONObject.toJavaObject(jsonObject0, AppEventMessage.class);
// appEventMessage.setId(nDid);
// appEventMessageTemplate.sendMember(appEventMessage);
// break;
// //实时数据
// case 1:
// JSONObject jsonObject2 = JSONObject.parseObject(JSON.toJSONString(dataDto));
// AppAutoDataMessage appAutoDataMessage = JSONObject.toJavaObject(jsonObject2, AppAutoDataMessage.class);
// appAutoDataMessage.setId(nDid);
// rtFeignClient.analysis(appAutoDataMessage);
// break;
// //处理主动上送的统计数据、电度数据
// case 2:
// case 3:
// JSONObject jsonObject3 = JSONObject.parseObject(JSON.toJSONString(dataDto));
// AppAutoDataMessage appAutoDataMessage2 = JSONObject.toJavaObject(jsonObject3, AppAutoDataMessage.class);
// appAutoDataMessage2.setId(nDid);
// appAutoDataMessageTemplate.sendMember(appAutoDataMessage2);
// break;
// }
// break;
// default:
// break;
// }
// }
@MqttSubscribe(value = "/Pfm/DevFileRsp/{version}/{edgeId}",qos = 1) @MqttSubscribe(value = "/Pfm/DevFileRsp/{version}/{edgeId}",qos = 1)
public void file(String topic, MqttMessage message, @NamedValue("version") String version, @NamedValue("edgeId") String nDid, @Payload String payload) { public void file(String topic, MqttMessage message, @NamedValue("version") String version, @NamedValue("edgeId") String nDid, @Payload String payload) {
Gson gson = new Gson(); Gson gson = new Gson();
@@ -1177,7 +725,7 @@ public class MqttMessageHandler {
log.info("装置根目录应答"); log.info("装置根目录应答");
redisUtil.saveByKeyWithExpire(AppRedisKey.DEVICE_ROOT_PATH + nDid,fileDto.getMsg().getName(),10L); redisUtil.saveByKeyWithExpire(AppRedisKey.DEVICE_ROOT_PATH + nDid,fileDto.getMsg().getName(),10L);
case 4664: case 4664:
log.info("上送程序文件单帧,设备应答"); //log.info("上送程序文件单帧,设备应答");
FileRedisDto dto2 = new FileRedisDto(); FileRedisDto dto2 = new FileRedisDto();
dto2.setCode(fileDto.getCode()); dto2.setCode(fileDto.getCode());
redisUtil.saveByKeyWithExpire(AppRedisKey.UPLOAD.concat(nDid).concat(String.valueOf(fileDto.getMid())),dto2,10L); redisUtil.saveByKeyWithExpire(AppRedisKey.UPLOAD.concat(nDid).concat(String.valueOf(fileDto.getMid())),dto2,10L);
@@ -1188,28 +736,39 @@ public class MqttMessageHandler {
&& !Objects.isNull(fileDto.getMsg().getFileCrc())) { && !Objects.isNull(fileDto.getMsg().getFileCrc())) {
CsEdDataPO po = csEdDataFeignClient.findByPath(fileDto.getMsg().getFileName()).getData(); CsEdDataPO po = csEdDataFeignClient.findByPath(fileDto.getMsg().getFileName()).getData();
if (Objects.nonNull(po)) { if (Objects.nonNull(po)) {
ReqAndResDto.Req pojo;
Object object = channelObjectUtil.getDeviceMid(nDid);
if (!Objects.isNull(object)) {
mid = (Integer) object;
}
if (Objects.equals(fileDto.getMsg().getFileCrc(), po.getCrc())) { if (Objects.equals(fileDto.getMsg().getFileCrc(), po.getCrc())) {
log.info("设备收到程序文件应答,开始升级"); log.info("设备收到程序文件应答,开始升级");
pojo = getPojo(mid,1); //pojo = getPojo(fileDto.getMid(),1);
String json = "{type:crcCheck,systemCrc:" + po.getCrc()
+ ",devCrc:" + fileDto.getMsg().getFileCrc()
+ ",mid:" + fileDto.getMid()
+ ",version:" + version
+ "}";
publisher.send("/Web/Progress/UpgradeFile/" + nDid, new Gson().toJson(json), 1, false);
} else { } else {
log.info("设备收到程序文件应答,文件校验失败"); log.info("设备收到程序文件应答,文件校验失败");
pojo = getPojo(mid,0); ReqAndResDto.Req pojo = getPojo(fileDto.getMid(),0);
}
publisher.send("/Pfm/DevFileCmd/" + version + "/" + nDid, new Gson().toJson(pojo), 1, false); publisher.send("/Pfm/DevFileCmd/" + version + "/" + nDid, new Gson().toJson(pojo), 1, false);
mid = mid + 1;
if (mid > 10000) {
mid = 1;
}
redisUtil.saveByKey(AppRedisKey.DEVICE_MID + nDid,mid);
} }
} }
if (Objects.equals(fileDto.getCode(),200) && Objects.isNull(fileDto.getMsg().getCmpAlg())) { }
log.info("装置启动升级,回复平台升级结果"); if (Objects.isNull(fileDto.getMsg())) {
CsEquipmentDeliveryPO po = equipmentFeignClient.findDevByNDid(nDid).getData();
log.info("装置升级,回复平台升级结果,设备是:{}", nDid);
if (Objects.equals(fileDto.getCode(),200)) {
Object object = redisUtil.getObjectByKey("chk:" + nDid);
if (Objects.nonNull(object) && (Integer)object == 1) {
csUpgradeLogsFeignClient.update(po.getId(), String.valueOf(fileDto.getCode()),1);
//重启设备
askDeviceDataFeignClient.rebootDevice(nDid);
}
//校验码没有匹配上,升级失败的
else {
csUpgradeLogsFeignClient.update(po.getId(), "0",0);
}
} else {
csUpgradeLogsFeignClient.update(po.getId(), String.valueOf(fileDto.getCode()),0);
}
} }
default: default:
break; break;
@@ -1227,95 +786,18 @@ public class MqttMessageHandler {
UpgradeDevDto upgradeDevDto = new UpgradeDevDto(); UpgradeDevDto upgradeDevDto = new UpgradeDevDto();
upgradeDevDto.setChk(chk); upgradeDevDto.setChk(chk);
reqAndResParam.setMsg(upgradeDevDto);
return reqAndResParam; return reqAndResParam;
} }
// /** /**
// * 文件传输 * 装置异常事件记录
// * @param topic * @param topic
// * @param message * @param message
// * @param version * @param version
// * @param nDid * @param nDid
// * @param payload * @param payload
// */ */
// @MqttSubscribe(value = "/Pfm/DevFileRsp/{version}/{edgeId}",qos = 1)
// @Transactional(rollbackFor = Exception.class)
// public void file(String topic, MqttMessage message, @NamedValue("version") String version, @NamedValue("edgeId") String nDid, @Payload String payload) {
// //解析数据
// Gson gson = new Gson();
// FileDto fileDto = gson.fromJson(new String(message.getPayload(), StandardCharsets.UTF_8), FileDto.class);
// JSONObject jsonObject = JSONObject.parseObject(JSON.toJSONString(fileDto));
// AppFileMessage appFileMessage = JSONObject.toJavaObject(jsonObject, AppFileMessage.class);
// appFileMessage.setId(nDid);
// //响应请求
// switch (fileDto.getType()){
// case 4657:
// if (Objects.equals(fileDto.getCode(),AccessEnum.SUCCESS.getCode())) {
// String key = AppRedisKey.PROJECT_INFO + nDid;
// if (Objects.isNull(fileDto.getMsg().getType())) {
// handleDefaultCase(fileDto, nDid);
// } else {
// if (Objects.equals("dir", fileDto.getMsg().getType())) {
// saveDirectoryInfo(fileDto.getMsg().getDirInfo(), key);
// } else if (Objects.equals("file", fileDto.getMsg().getType())){
// saveFileInfo(fileDto.getMsg().getFileInfo(), key);
// appFileMessageTemplate.sendMember(appFileMessage);
// }
// }
// } else if (Objects.equals(fileDto.getCode(),AccessEnum.NOT_FIND.getCode())) {
// Object object = redisUtil.getObjectByKey("fileMid:" + nDid);
// if (Objects.nonNull(object)) {
// String data = redisUtil.getObjectByKey("fileMid:" + nDid).toString();
// String [] arr = data.split("concat");
// Integer mid = Integer.parseInt(arr[0]);
// String fileName = arr[1];
// if (Objects.equals(mid,fileDto.getMid())) {
// List<WaveTimeDto> list = channelObjectUtil.objectToList( redisUtil.getObjectByKey("eventFile:" + nDid),WaveTimeDto.class);
// list.removeIf(item -> item.getFileName().equals(fileName));
// redisUtil.saveByKey("eventFile:" + nDid, list);
// if (CollectionUtil.isNotEmpty(list)) {
// redisUtil.delete("handleEvent:" + nDid);
// waveFeignClient.channelWave(nDid);
// }
// }
// }
// }
// break;
// case 4658:
// FileRedisDto dto = new FileRedisDto();
// dto.setCode(fileDto.getCode());
// redisUtil.saveByKeyWithExpire(AppRedisKey.DOWNLOAD + fileDto.getMsg().getName() + fileDto.getMid(),dto,60L);
// if (Objects.equals(fileDto.getCode(),AccessEnum.SUCCESS.getCode())){
// appFileStreamMessageTemplate.sendMember(appFileMessage);
// }
// //todo 处理文件信息,先缓存起来,后期在询问
// else if (Objects.equals(fileDto.getCode(),AccessEnum.REFUSE_WAIT.getCode())) {
// log.info("需要缓存请求的文件信息");
// }
// break;
// case 4659:
// log.info("装置收到系统上传的文件");
// FileRedisDto fileRedisDto = new FileRedisDto();
// fileRedisDto.setCode(fileDto.getCode());
// redisUtil.saveByKeyWithExpire(AppRedisKey.UPLOAD.concat(nDid).concat(String.valueOf(fileDto.getMid())),fileRedisDto,10L);
// redisUtil.saveByKeyWithExpire("uploading","uploading",20L);
// break;
// case 4660:
// log.info("设备目录/文件删除应答");
// redisUtil.saveByKeyWithExpire( "deleteDir"+ nDid,fileDto.getCode(),10L);
// break;
// case 4661:
// log.info("设备目录创建应答");
// redisUtil.saveByKeyWithExpire( "createDir"+ nDid,fileDto.getCode(),10L);
// break;
// case 4662:
// log.info("装置根目录应答");
// redisUtil.saveByKeyWithExpire(AppRedisKey.DEVICE_ROOT_PATH + nDid,fileDto.getMsg().getName(),10L);
// default:
// break;
// }
// }
@MqttSubscribe(value = "/Dev/Error/{edgeId}",qos = 1) @MqttSubscribe(value = "/Dev/Error/{edgeId}",qos = 1)
public void devErrorInfo(String topic, MqttMessage message, @NamedValue("version") String version, @NamedValue("edgeId") String nDid, @Payload String payload) { public void devErrorInfo(String topic, MqttMessage message, @NamedValue("version") String version, @NamedValue("edgeId") String nDid, @Payload String payload) {
Gson gson = new Gson(); Gson gson = new Gson();
@@ -1332,24 +814,6 @@ public class MqttMessageHandler {
} }
}); });
} }
// /**
// * 装置异常事件记录
// * @param topic
// * @param message
// * @param version
// * @param nDid
// * @param payload
// */
// @MqttSubscribe(value = "/Dev/Error/{edgeId}",qos = 1)
// public void devErrorInfo(String topic, MqttMessage message, @NamedValue("version") String version, @NamedValue("edgeId") String nDid, @Payload String payload) {
// //解析数据
// Gson gson = new Gson();
// EventDto eventDto = gson.fromJson(new String(message.getPayload(), StandardCharsets.UTF_8), EventDto.class);
// JSONObject jsonObject0 = JSONObject.parseObject(JSON.toJSONString(eventDto));
// AppEventMessage appEventMessage = JSONObject.toJavaObject(jsonObject0, AppEventMessage.class);
// appEventMessage.setId(nDid);
// appEventMessageTemplate.sendMember(appEventMessage);
// }
private void saveDirectoryInfo(List<FileDto.DirInfo> dirInfo, String key) { private void saveDirectoryInfo(List<FileDto.DirInfo> dirInfo, String key) {
if (!CollectionUtil.isEmpty(dirInfo)) { if (!CollectionUtil.isEmpty(dirInfo)) {