2026-06-24 f4bd1f3c89d906131495a0aca5aaf82966378510
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
package cn.iocoder.yudao.module.iot.service.ota;
 
import cn.hutool.core.collection.CollUtil;
import cn.hutool.core.convert.Convert;
import cn.hutool.core.lang.Assert;
import cn.hutool.core.util.StrUtil;
import cn.iocoder.yudao.framework.common.pojo.PageResult;
import cn.iocoder.yudao.framework.common.util.json.JsonUtils;
import cn.iocoder.yudao.framework.common.util.object.BeanUtils;
import cn.iocoder.yudao.module.iot.controller.admin.ota.vo.task.record.IotOtaTaskRecordPageReqVO;
import cn.iocoder.yudao.module.iot.core.enums.IotDeviceMessageMethodEnum;
import cn.iocoder.yudao.module.iot.core.mq.message.IotDeviceMessage;
import cn.iocoder.yudao.module.iot.core.topic.ota.IotDeviceOtaProgressReqDTO;
import cn.iocoder.yudao.module.iot.core.topic.ota.IotDeviceOtaUpgradeReqDTO;
import cn.iocoder.yudao.module.iot.dal.dataobject.device.IotDeviceDO;
import cn.iocoder.yudao.module.iot.dal.dataobject.ota.IotOtaFirmwareDO;
import cn.iocoder.yudao.module.iot.dal.dataobject.ota.IotOtaTaskRecordDO;
import cn.iocoder.yudao.module.iot.dal.mysql.ota.IotOtaTaskRecordMapper;
import cn.iocoder.yudao.module.iot.enums.ota.IotOtaTaskRecordStatusEnum;
import cn.iocoder.yudao.module.iot.service.device.IotDeviceService;
import cn.iocoder.yudao.module.iot.service.device.message.IotDeviceMessageService;
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.util.*;
 
import static cn.iocoder.yudao.framework.common.exception.util.ServiceExceptionUtil.exception;
import static cn.iocoder.yudao.framework.common.util.collection.CollectionUtils.*;
import static cn.iocoder.yudao.module.iot.enums.ErrorCodeConstants.*;
 
/**
 * OTA 升级任务记录 Service 实现类
 */
@Service
@Validated
@Slf4j
public class IotOtaTaskRecordServiceImpl implements IotOtaTaskRecordService {
 
    @Resource
    private IotOtaTaskRecordMapper otaTaskRecordMapper;
 
    @Resource
    private IotOtaFirmwareService otaFirmwareService;
    @Resource
    private IotOtaTaskService otaTaskService;
    @Resource
    private IotDeviceMessageService deviceMessageService;
    @Resource
    private IotDeviceService deviceService;
 
    @Override
    public void createOtaTaskRecordList(List<IotDeviceDO> devices, Long firmwareId, Long taskId) {
        List<IotOtaTaskRecordDO> records = convertList(devices, device ->
                IotOtaTaskRecordDO.builder().firmwareId(firmwareId).taskId(taskId)
                        .deviceId(device.getId()).fromFirmwareId(Convert.toLong(device.getFirmwareId()))
                        .status(IotOtaTaskRecordStatusEnum.PENDING.getStatus()).progress(0).build());
        otaTaskRecordMapper.insertBatch(records);
    }
 
    @Override
    public Map<Integer, Long> getOtaTaskRecordStatusStatistics(Long firmwareId, Long taskId) {
        // 按照 status 枚举,初始化 countMap 为 0
        Map<Integer, Long> countMap = convertMap(Arrays.asList(IotOtaTaskRecordStatusEnum.values()),
                IotOtaTaskRecordStatusEnum::getStatus, iotOtaTaskRecordStatusEnum -> 0L);
 
        // 查询记录,只返回 id、status 字段
        List<IotOtaTaskRecordDO> records = otaTaskRecordMapper.selectListByFirmwareIdAndTaskId(firmwareId, taskId);
        Map<Long, List<Integer>> deviceStatusesMap = convertMultiMap(records,
                IotOtaTaskRecordDO::getDeviceId, IotOtaTaskRecordDO::getStatus);
        // 找到第一个匹配的优先级状态,避免重复计算
        deviceStatusesMap.forEach((deviceId, statuses) -> {
            for (Integer priorityStatus : IotOtaTaskRecordStatusEnum.PRIORITY_STATUSES) {
                if (statuses.contains(priorityStatus)) {
                    countMap.put(priorityStatus, countMap.get(priorityStatus) + 1);
                    return;
                }
            }
        });
        return countMap;
    }
 
    @Override
    public IotOtaTaskRecordDO getOtaTaskRecord(Long id) {
        return otaTaskRecordMapper.selectById(id);
    }
 
    @Override
    public PageResult<IotOtaTaskRecordDO> getOtaTaskRecordPage(IotOtaTaskRecordPageReqVO pageReqVO) {
        return otaTaskRecordMapper.selectPage(pageReqVO);
    }
 
    @Override
    public void cancelTaskRecordListByTaskId(Long taskId) {
        List<IotOtaTaskRecordDO> records = otaTaskRecordMapper.selectListByTaskIdAndStatus(
                taskId, IotOtaTaskRecordStatusEnum.IN_PROCESS_STATUSES);
        if (CollUtil.isEmpty(records)) {
            return;
        }
        // 批量更新
        Collection<Long> ids = convertSet(records, IotOtaTaskRecordDO::getId);
        otaTaskRecordMapper.updateListByIdAndStatus(ids, IotOtaTaskRecordStatusEnum.IN_PROCESS_STATUSES,
                IotOtaTaskRecordDO.builder().status(IotOtaTaskRecordStatusEnum.CANCELED.getStatus())
                        .description(IotOtaTaskRecordDO.DESCRIPTION_CANCEL_BY_TASK).build());
    }
 
    @Override
    public List<IotOtaTaskRecordDO> getOtaTaskRecordListByDeviceIdAndStatus(Set<Long> deviceIds, Set<Integer> statuses) {
        return otaTaskRecordMapper.selectListByDeviceIdAndStatus(deviceIds, statuses);
    }
 
    @Override
    public List<IotOtaTaskRecordDO> getOtaRecordListByStatus(Integer status) {
        return otaTaskRecordMapper.selectListByStatus(status);
    }
 
    @Override
    public void cancelOtaTaskRecord(Long id) {
        // 1. 校验记录是否存在
        IotOtaTaskRecordDO record = validateUpgradeRecordExists(id);
 
        // 2. 更新记录状态为取消
        int updateCount = otaTaskRecordMapper.updateByIdAndStatus(record.getId(), IotOtaTaskRecordStatusEnum.IN_PROCESS_STATUSES,
                IotOtaTaskRecordDO.builder().id(id).status(IotOtaTaskRecordStatusEnum.CANCELED.getStatus())
                .description(IotOtaTaskRecordDO.DESCRIPTION_CANCEL_BY_RECORD).build());
        if (updateCount == 0) {
            throw exception(OTA_TASK_RECORD_CANCEL_FAIL_STATUS_ERROR);
        }
 
        // 3. 检查并更新任务状态
        checkAndUpdateOtaTaskStatus(record.getTaskId());
    }
 
    @Override
    public boolean pushOtaTaskRecord(IotOtaTaskRecordDO record, IotOtaFirmwareDO fireware, IotDeviceDO device) {
        try {
            // 1. 推送 OTA 任务记录
            IotDeviceOtaUpgradeReqDTO params = BeanUtils.toBean(fireware, IotDeviceOtaUpgradeReqDTO.class);
            IotDeviceMessage message = IotDeviceMessage.requestOf(
                    IotDeviceMessageMethodEnum.OTA_UPGRADE.getMethod(), params);
            deviceMessageService.sendDeviceMessage(message, device);
 
            // 2. 更新 OTA 升级记录状态为进行中
            int updateCount = otaTaskRecordMapper.updateByIdAndStatus(
                    record.getId(), IotOtaTaskRecordStatusEnum.PENDING.getStatus(),
                    IotOtaTaskRecordDO.builder().status(IotOtaTaskRecordStatusEnum.PUSHED.getStatus())
                            .description(StrUtil.format("已推送,设备消息编号({})", message.getId())).build());
            Assert.isTrue(updateCount == 1, "更新设备记录({})状态失败", record.getId());
            return true;
        } catch (Exception ex) {
            log.error("[pushOtaTaskRecord][推送 OTA 任务记录({}) 失败]", record.getId(), ex);
            otaTaskRecordMapper.updateById(IotOtaTaskRecordDO.builder().id(record.getId())
                    .description(StrUtil.format("推送失败,错误信息({})", ex.getMessage())).build());
            return false;
        }
    }
 
    private IotOtaTaskRecordDO validateUpgradeRecordExists(Long id) {
        IotOtaTaskRecordDO upgradeRecord = otaTaskRecordMapper.selectById(id);
        if (upgradeRecord == null) {
            throw exception(OTA_TASK_RECORD_NOT_EXISTS);
        }
        return upgradeRecord;
    }
 
    @Override
    @Transactional(rollbackFor = Exception.class)
    public void updateOtaRecordProgress(IotDeviceDO device, IotDeviceMessage message) {
        // 1.1 参数解析
        IotDeviceOtaProgressReqDTO params = JsonUtils.convertObject(message.getParams(), IotDeviceOtaProgressReqDTO.class);
        String version = params.getVersion();
        Assert.notBlank(version, "version 不能为空");
        Integer status = params.getStatus();
        Assert.notNull(status, "status 不能为空");
        Assert.notNull(IotOtaTaskRecordStatusEnum.of(status), "status 状态不正确");
        String description = params.getDescription();
        Integer progress = params.getProgress();
        Assert.notNull(progress, "progress 不能为空");
        Assert.isTrue(progress >= 0 && progress <= 100, "progress 必须在 0-100 之间");
        // 1.2 查询 OTA 升级记录
        List<IotOtaTaskRecordDO> records = otaTaskRecordMapper.selectListByDeviceIdAndStatus(
                device.getId(), IotOtaTaskRecordStatusEnum.IN_PROCESS_STATUSES);
        if (CollUtil.isEmpty(records)) {
            throw exception(OTA_TASK_RECORD_UPDATE_PROGRESS_FAIL_NO_EXISTS);
        }
        if (records.size() > 1) {
            log.warn("[updateOtaRecordProgress][message({}) 对应升级记录过多({})]", message, records);
        }
        IotOtaTaskRecordDO record = CollUtil.getFirst(records);
        // 1.3 查询 OTA 固件
        IotOtaFirmwareDO firmware = otaFirmwareService.getOtaFirmwareByProductIdAndVersion(
                device.getProductId(), version);
        if (firmware == null) {
            throw exception(OTA_FIRMWARE_NOT_EXISTS);
        }
 
        // 2. 更新 OTA 升级记录状态
        int updateCount = otaTaskRecordMapper.updateByIdAndStatus(
                record.getId(), IotOtaTaskRecordStatusEnum.IN_PROCESS_STATUSES,
                IotOtaTaskRecordDO.builder().status(status).description(description).progress(progress).build());
        if (updateCount == 0) {
            throw exception(OTA_TASK_RECORD_UPDATE_PROGRESS_FAIL_NO_EXISTS);
        }
 
        // 3. 如果升级成功,则更新设备固件版本
        if (IotOtaTaskRecordStatusEnum.SUCCESS.getStatus().equals(status)) {
            deviceService.updateDeviceFirmware(device.getId(), firmware.getId());
        }
 
        // 4. 如果状态是“已结束”(非进行中),则更新任务状态
        if (!IotOtaTaskRecordStatusEnum.IN_PROCESS_STATUSES.contains(status)) {
            checkAndUpdateOtaTaskStatus(record.getTaskId());
        }
    }
 
    /**
     * 检查并更新任务状态
     * 如果任务下没有进行中的记录,则将任务状态更新为已结束
     */
    private void checkAndUpdateOtaTaskStatus(Long taskId) {
        // 如果还有进行中的记录,直接返回
        Long inProcessCount = otaTaskRecordMapper.selectCountByTaskIdAndStatus(
                taskId, IotOtaTaskRecordStatusEnum.IN_PROCESS_STATUSES);
        if (inProcessCount > 0) {
            return;
        }
 
        // 没有进行中的记录,将任务状态更新为已结束
        otaTaskService.updateOtaTaskStatusEnd(taskId);
    }
 
}