From 20e5bb3e933db3e23761c82204340e17ec198d9c Mon Sep 17 00:00:00 2001 From: xy <748613696@qq.com> Date: Wed, 22 Jul 2026 08:54:08 +0800 Subject: [PATCH] =?UTF-8?q?feat(mqtt):=20=E6=B7=BB=E5=8A=A0=E8=AE=BE?= =?UTF-8?q?=E5=A4=87=E7=A8=8B=E5=BA=8F=E6=96=87=E4=BB=B6=E4=B8=8A=E4=BC=A0?= =?UTF-8?q?=E5=92=8C=E5=8D=87=E7=BA=A7=E5=8A=9F=E8=83=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 在FileDto中新增压缩方式、文件大小、文件个数等文件传输相关字段 - 在MqttMessageHandler中添加4664消息类型处理设备程序文件单帧上送 - 实现文件校验码验证和程序升级请求逻辑 - 添加UpgradeDevDto用于程序升级相关的数据传输 - 在CsDevModelServiceImpl异常处理中清理上传状态缓存 - 添加程序升级(type 32)消息类型定义 - 集成csEdDataFeignClient进行设备数据查询 - 实现升级结果确认和设备端升级流程控制 --- .../java/com/njcn/access/enums/TypeEnum.java | 2 + .../njcn/access/pojo/dto/UpgradeDevDto.java | 48 +++++++++++++++++ .../njcn/access/pojo/dto/file/FileDto.java | 24 +++++++++ .../access/handler/MqttMessageHandler.java | 52 ++++++++++++++++++- .../service/impl/CsDevModelServiceImpl.java | 2 + .../zlevent/service/impl/FileServiceImpl.java | 2 + 6 files changed, 129 insertions(+), 1 deletion(-) create mode 100644 iot-access/access-api/src/main/java/com/njcn/access/pojo/dto/UpgradeDevDto.java diff --git a/iot-access/access-api/src/main/java/com/njcn/access/enums/TypeEnum.java b/iot-access/access-api/src/main/java/com/njcn/access/enums/TypeEnum.java index b7db8b1..6f7b98a 100644 --- a/iot-access/access-api/src/main/java/com/njcn/access/enums/TypeEnum.java +++ b/iot-access/access-api/src/main/java/com/njcn/access/enums/TypeEnum.java @@ -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","设备重启"), diff --git a/iot-access/access-api/src/main/java/com/njcn/access/pojo/dto/UpgradeDevDto.java b/iot-access/access-api/src/main/java/com/njcn/access/pojo/dto/UpgradeDevDto.java new file mode 100644 index 0000000..bde0d29 --- /dev/null +++ b/iot-access/access-api/src/main/java/com/njcn/access/pojo/dto/UpgradeDevDto.java @@ -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-zip,2-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; +} diff --git a/iot-access/access-api/src/main/java/com/njcn/access/pojo/dto/file/FileDto.java b/iot-access/access-api/src/main/java/com/njcn/access/pojo/dto/file/FileDto.java index 2a57019..6817eef 100644 --- a/iot-access/access-api/src/main/java/com/njcn/access/pojo/dto/file/FileDto.java +++ b/iot-access/access-api/src/main/java/com/njcn/access/pojo/dto/file/FileDto.java @@ -69,6 +69,30 @@ public class FileDto implements Serializable { @SerializedName("Offset") private Integer offset; + @SerializedName("CmpAlg") + @ApiModelProperty("压缩方式(0-无,1-zip,2-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 diff --git a/iot-access/access-boot/src/main/java/com/njcn/access/handler/MqttMessageHandler.java b/iot-access/access-boot/src/main/java/com/njcn/access/handler/MqttMessageHandler.java index 7e4d064..e597554 100644 --- a/iot-access/access-boot/src/main/java/com/njcn/access/handler/MqttMessageHandler.java +++ b/iot-access/access-boot/src/main/java/com/njcn/access/handler/MqttMessageHandler.java @@ -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 diff --git a/iot-access/access-boot/src/main/java/com/njcn/access/service/impl/CsDevModelServiceImpl.java b/iot-access/access-boot/src/main/java/com/njcn/access/service/impl/CsDevModelServiceImpl.java index cb97631..ae2a96a 100644 --- a/iot-access/access-boot/src/main/java/com/njcn/access/service/impl/CsDevModelServiceImpl.java +++ b/iot-access/access-boot/src/main/java/com/njcn/access/service/impl/CsDevModelServiceImpl.java @@ -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("系统上传文件"); diff --git a/iot-analysis/analysis-zl-event/zl-event-boot/src/main/java/com/njcn/zlevent/service/impl/FileServiceImpl.java b/iot-analysis/analysis-zl-event/zl-event-boot/src/main/java/com/njcn/zlevent/service/impl/FileServiceImpl.java index 64d882d..95034d8 100644 --- a/iot-analysis/analysis-zl-event/zl-event-boot/src/main/java/com/njcn/zlevent/service/impl/FileServiceImpl.java +++ b/iot-analysis/analysis-zl-event/zl-event-boot/src/main/java/com/njcn/zlevent/service/impl/FileServiceImpl.java @@ -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;