feat(mqtt): 添加设备程序文件上传和升级功能

- 在FileDto中新增压缩方式、文件大小、文件个数等文件传输相关字段
- 在MqttMessageHandler中添加4664消息类型处理设备程序文件单帧上送
- 实现文件校验码验证和程序升级请求逻辑
- 添加UpgradeDevDto用于程序升级相关的数据传输
- 在CsDevModelServiceImpl异常处理中清理上传状态缓存
- 添加程序升级(type 32)消息类型定义
- 集成csEdDataFeignClient进行设备数据查询
- 实现升级结果确认和设备端升级流程控制
This commit is contained in:
xy
2026-07-22 08:54:08 +08:00
parent 759a811247
commit 20e5bb3e93
6 changed files with 129 additions and 1 deletions

View File

@@ -45,11 +45,13 @@ public enum TypeEnum {
TYPE_29("9217","设备心跳请求"),
TYPE_30("4865","设备数据主动上送"),
TYPE_31("8503","设备控制命令"),
TYPE_32("8504","程序升级"),
READ_FILE_DIR("1101", "读取文件目录"),
FILE_DOWNLOAD("1102", "文件下载"),
FIXED_VALUE("1103", "定值读取/写入"),
INNER_FIXED_VALUE("1104", "内部定值读取/写入"),
WORKING_LOG("1111","设备运行日志"),
DEVICE_VERSION("1112","设备版本信息"),
DEVICE_REBOOT("1114","设备重启"),

View File

@@ -0,0 +1,48 @@
package com.njcn.access.pojo.dto;
import com.alibaba.nacos.shaded.com.google.gson.annotations.SerializedName;
import io.swagger.annotations.ApiModelProperty;
import lombok.Data;
/**
* @author xy
*/
@Data
public class UpgradeDevDto {
@SerializedName("CmpAlg")
@ApiModelProperty("压缩方式0-无1-zip2-tar")
private Integer cmpAlg;
@SerializedName("RawSize")
@ApiModelProperty("压缩前全部文件总大小(单位字节)")
private Integer rawSize;
@SerializedName("FileCnt")
@ApiModelProperty("压缩前全部文件个数")
private Integer fileCnt;
@SerializedName("Name")
@ApiModelProperty("文件名称,全路径")
private String name;
@SerializedName("FileSize")
@ApiModelProperty("文件总大小(单位字节)")
private Integer fileSize;
@SerializedName("Offset")
@ApiModelProperty("当前上送数据包在文件中偏移位置")
private Integer offset;
@SerializedName("Len")
@ApiModelProperty("当前上送数据包长度")
private Integer len;
@SerializedName("Data")
@ApiModelProperty("文件包数据")
private String data;
@SerializedName("Chk")
@ApiModelProperty("确认升级0-取消升级1-确认升级)")
private Integer chk;
}

View File

@@ -69,6 +69,30 @@ public class FileDto implements Serializable {
@SerializedName("Offset")
private Integer offset;
@SerializedName("CmpAlg")
@ApiModelProperty("压缩方式0-无1-zip2-tar")
private Integer cmpAlg;
@SerializedName("RawSize")
@ApiModelProperty("压缩前全部文件总大小(单位字节)")
private Integer rawSize;
@SerializedName("FileCnt")
@ApiModelProperty("压缩前全部文件个数")
private Integer fileCnt;
@SerializedName("Len")
@ApiModelProperty("当前上送数据包长度大小限定每包不超过50KB")
private Integer len;
@SerializedName("FileName")
@ApiModelProperty("文件名")
private String fileName;
@SerializedName("FileCrc")
@ApiModelProperty("文件校验码")
private String fileCrc;
}
@Data

View File

@@ -51,7 +51,6 @@ import org.eclipse.paho.client.mqttv3.MqttMessage;
import org.springframework.beans.BeanUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import javax.validation.ConstraintViolation;
import javax.validation.Validator;
@@ -102,6 +101,8 @@ public class MqttMessageHandler {
private final LogMessageTemplate logMessageTemplate;
private final CsDeviceRegistryFeignClient csDeviceRegistryFeignClient;
private final ExecutorService mqttMessageExecutor;
private final CsEdDataFeignClient csEdDataFeignClient;
private static Integer mid = 1;
@Autowired
Validator validator;
@@ -1175,11 +1176,60 @@ public class MqttMessageHandler {
case 4662:
log.info("装置根目录应答");
redisUtil.saveByKeyWithExpire(AppRedisKey.DEVICE_ROOT_PATH + nDid,fileDto.getMsg().getName(),10L);
case 4664:
log.info("上送程序文件单帧,设备应答");
FileRedisDto dto2 = new FileRedisDto();
dto2.setCode(fileDto.getCode());
redisUtil.saveByKeyWithExpire(AppRedisKey.UPLOAD.concat(nDid).concat(String.valueOf(fileDto.getMid())),dto2,10L);
redisUtil.saveByKeyWithExpire("uploading:" + nDid,"uploading",20L);
//判断是否校验完成,程序升级请求
if (Objects.equals(fileDto.getCode(),201)
&& !Objects.isNull(fileDto.getMsg().getFileName())
&& !Objects.isNull(fileDto.getMsg().getFileCrc())) {
CsEdDataPO po = csEdDataFeignClient.findByPath(fileDto.getMsg().getFileName()).getData();
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())) {
log.info("设备收到程序文件应答,开始升级");
pojo = getPojo(mid,1);
} else {
log.info("设备收到程序文件应答,文件校验失败");
pojo = getPojo(mid,0);
}
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("装置启动升级,回复平台升级结果");
}
default:
break;
}
}
public ReqAndResDto.Req getPojo(Integer mid, Integer chk) {
//组装报文
ReqAndResDto.Req reqAndResParam = new ReqAndResDto.Req();
reqAndResParam.setMid(mid);
reqAndResParam.setDid(0);
reqAndResParam.setPri(AccessEnum.FIRST_CHANNEL.getCode());
reqAndResParam.setType(Integer.parseInt(TypeEnum.TYPE_32.getCode()));
reqAndResParam.setExpire(-1);
UpgradeDevDto upgradeDevDto = new UpgradeDevDto();
upgradeDevDto.setChk(chk);
return reqAndResParam;
}
// /**
// * 文件传输
// * @param topic

View File

@@ -251,6 +251,8 @@ public class CsDevModelServiceImpl implements ICsDevModelService {
sendNextStep(logDto,path,file,length,bytes,0,version,id,1,hexString,false);
}
} catch (Exception e) {
redisUtil.delete("uploading");
redisUtil.delete("fileDowning:" + id);
LogMessage logDto = new LogMessage();
logDto.setResult(0);
logDto.setOperate("系统上传文件");

View File

@@ -12,6 +12,8 @@ import com.njcn.access.enums.AccessEnum;
import com.njcn.access.enums.AccessResponseEnum;
import com.njcn.access.enums.TypeEnum;
import com.njcn.access.pojo.dto.ReqAndResDto;
import com.njcn.zlevent.pojo.dto.FileDownloadRequestDTO;
import com.njcn.zlevent.pojo.dto.FileDownloadResponeDTO;
import com.njcn.access.pojo.dto.file.FileDto;
import com.njcn.access.utils.*;
import com.njcn.advance.api.EventCauseFeignClient;