2026-06-30 24681c81c09022f584a57006f2534b5f74723414
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
package cn.iocoder.yudao.module.iot.gateway.protocol.modbus.tcpserver.manager;
 
import cn.hutool.core.collection.CollUtil;
import cn.iocoder.yudao.framework.common.enums.CommonStatusEnum;
import cn.iocoder.yudao.framework.common.pojo.CommonResult;
import cn.iocoder.yudao.module.iot.core.biz.IotDeviceCommonApi;
import cn.iocoder.yudao.module.iot.core.biz.dto.IotModbusDeviceConfigListReqDTO;
import cn.iocoder.yudao.module.iot.core.biz.dto.IotModbusDeviceConfigRespDTO;
import cn.iocoder.yudao.module.iot.core.enums.modbus.IotModbusModeEnum;
import cn.iocoder.yudao.module.iot.core.enums.IotProtocolTypeEnum;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
 
import java.util.Collections;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
 
import static cn.iocoder.yudao.framework.common.util.collection.CollectionUtils.convertSet;
 
/**
 * IoT Modbus TCP Server 配置缓存:认证时按需加载,断连时清理,定时刷新已连接设备
 *
 * @author 芋道源码
 */
@RequiredArgsConstructor
@Slf4j
public class IotModbusTcpServerConfigCacheService {
 
    private final IotDeviceCommonApi deviceApi;
 
    /**
     * 配置缓存:deviceId -> 配置
     */
    private final Map<Long, IotModbusDeviceConfigRespDTO> configCache = new ConcurrentHashMap<>();
 
    /**
     * 加载单个设备的配置(认证成功后调用)
     *
     * @param deviceId 设备 ID
     * @return 设备配置
     */
    public IotModbusDeviceConfigRespDTO loadDeviceConfig(Long deviceId) {
        try {
            // 1. 从远程 API 获取配置
            IotModbusDeviceConfigListReqDTO reqDTO = new IotModbusDeviceConfigListReqDTO()
                    .setStatus(CommonStatusEnum.ENABLE.getStatus())
                    .setMode(IotModbusModeEnum.POLLING.getMode())
                    .setProtocolType(IotProtocolTypeEnum.MODBUS_TCP_SERVER.getType())
                    .setDeviceIds(Collections.singleton(deviceId));
            CommonResult<List<IotModbusDeviceConfigRespDTO>> result = deviceApi.getModbusDeviceConfigList(reqDTO);
            result.checkError();
            IotModbusDeviceConfigRespDTO modbusConfig = CollUtil.getFirst(result.getData());
            if (modbusConfig == null) {
                log.warn("[loadDeviceConfig][远程获取配置失败,未找到数据, deviceId={}]", deviceId);
                return null;
            }
 
            // 2. 更新缓存并返回
            configCache.put(modbusConfig.getDeviceId(), modbusConfig);
            return modbusConfig;
        } catch (Exception e) {
            log.error("[loadDeviceConfig][从远程获取配置失败, deviceId={}]", deviceId, e);
            return null;
        }
    }
 
    /**
     * 刷新已连接设备的配置缓存
     * <p>
     * 定时调用,从远程 API 拉取最新配置,只更新已连接设备的缓存。
     *
     * @param connectedDeviceIds 当前已连接的设备 ID 集合
     * @return 已连接设备的最新配置列表
     */
    public List<IotModbusDeviceConfigRespDTO> refreshConnectedDeviceConfigList(Set<Long> connectedDeviceIds) {
        if (CollUtil.isEmpty(connectedDeviceIds)) {
            return Collections.emptyList();
        }
        try {
            // 1. 从远程获取已连接设备的配置
            CommonResult<List<IotModbusDeviceConfigRespDTO>> result = deviceApi.getModbusDeviceConfigList(
                    new IotModbusDeviceConfigListReqDTO().setStatus(CommonStatusEnum.ENABLE.getStatus())
                            .setMode(IotModbusModeEnum.POLLING.getMode())
                            .setProtocolType(IotProtocolTypeEnum.MODBUS_TCP_SERVER.getType())
                            .setDeviceIds(connectedDeviceIds));
            List<IotModbusDeviceConfigRespDTO> modbusConfigs = result.getCheckedData();
 
            // 2. 更新缓存并返回
            for (IotModbusDeviceConfigRespDTO config : modbusConfigs) {
                configCache.put(config.getDeviceId(), config);
            }
            return modbusConfigs;
        } catch (Exception e) {
            log.error("[refreshConnectedDeviceConfigList][刷新配置失败]", e);
            return null;
        }
    }
 
    /**
     * 清理本轮刷新后不再有效的设备配置
     *
     * @param refreshedDeviceIds 本轮参与刷新的设备编号
     * @param currentConfigs     本轮远端返回的有效配置
     * @return 本轮已不再有效的设备编号
     */
    public Set<Long> cleanupMissingConfigs(Set<Long> refreshedDeviceIds,
                                           List<IotModbusDeviceConfigRespDTO> currentConfigs) {
        if (CollUtil.isEmpty(refreshedDeviceIds)) {
            return Collections.emptySet();
        }
        Set<Long> currentDeviceIds = convertSet(currentConfigs, IotModbusDeviceConfigRespDTO::getDeviceId);
        Set<Long> missingDeviceIds = new HashSet<>(refreshedDeviceIds);
        missingDeviceIds.removeAll(currentDeviceIds);
        for (Long deviceId : missingDeviceIds) {
            configCache.remove(deviceId);
        }
        return missingDeviceIds;
    }
 
    /**
     * 获取设备配置
     */
    public IotModbusDeviceConfigRespDTO getConfig(Long deviceId) {
        IotModbusDeviceConfigRespDTO config = configCache.get(deviceId);
        if (config != null) {
            return config;
        }
        // 缓存未命中,从远程 API 获取
        return loadDeviceConfig(deviceId);
    }
 
    /**
     * 移除设备配置缓存(设备断连时调用)
     */
    public void removeConfig(Long deviceId) {
        configCache.remove(deviceId);
    }
 
}