package cn.iocoder.yudao.module.iot.gateway.protocol.modbus.tcpserver.handler.downstream; import cn.hutool.core.lang.Assert; import cn.hutool.core.util.ObjUtil; import cn.iocoder.yudao.module.iot.core.biz.dto.IotModbusDeviceConfigRespDTO; import cn.iocoder.yudao.module.iot.core.biz.dto.IotModbusPointRespDTO; import cn.iocoder.yudao.module.iot.core.enums.IotDeviceMessageMethodEnum; import cn.iocoder.yudao.module.iot.core.enums.modbus.IotModbusFrameFormatEnum; import cn.iocoder.yudao.module.iot.core.mq.message.IotDeviceMessage; import cn.iocoder.yudao.module.iot.gateway.protocol.modbus.common.utils.IotModbusCommonUtils; import cn.iocoder.yudao.module.iot.gateway.protocol.modbus.tcpserver.codec.IotModbusFrameEncoder; import cn.iocoder.yudao.module.iot.gateway.protocol.modbus.tcpserver.manager.IotModbusTcpServerConfigCacheService; import cn.iocoder.yudao.module.iot.gateway.protocol.modbus.tcpserver.manager.IotModbusTcpServerConnectionManager; import cn.iocoder.yudao.module.iot.gateway.protocol.modbus.tcpserver.manager.IotModbusTcpServerConnectionManager.ConnectionInfo; import lombok.extern.slf4j.Slf4j; import java.util.Map; import java.util.concurrent.atomic.AtomicInteger; /** * IoT Modbus TCP Server 下行消息处理器 *
* 负责:
* 1. 处理下行消息(如属性设置 thing.service.property.set)
* 2. 将属性值转换为 Modbus 写指令,通过 TCP 连接发送给设备
*
* @author 芋道源码
*/
@Slf4j
public class IotModbusTcpServerDownstreamHandler {
private final IotModbusTcpServerConnectionManager connectionManager;
private final IotModbusTcpServerConfigCacheService configCacheService;
private final IotModbusFrameEncoder frameEncoder;
/**
* TCP 事务 ID 自增器(与 PollScheduler 共享)
*/
private final AtomicInteger transactionIdCounter;
public IotModbusTcpServerDownstreamHandler(IotModbusTcpServerConnectionManager connectionManager,
IotModbusTcpServerConfigCacheService configCacheService,
IotModbusFrameEncoder frameEncoder,
AtomicInteger transactionIdCounter) {
this.connectionManager = connectionManager;
this.configCacheService = configCacheService;
this.frameEncoder = frameEncoder;
this.transactionIdCounter = transactionIdCounter;
}
/**
* 处理下行消息
*/
@SuppressWarnings({"unchecked", "DuplicatedCode"})
public void handle(IotDeviceMessage message) {
// 1.1 检查是否是属性设置消息
if (ObjUtil.equals(IotDeviceMessageMethodEnum.PROPERTY_POST.getMethod(), message.getMethod())) {
return;
}
if (ObjUtil.notEqual(IotDeviceMessageMethodEnum.PROPERTY_SET.getMethod(), message.getMethod())) {
log.debug("[handle][忽略非属性设置消息: {}]", message.getMethod());
return;
}
// 1.2 获取设备配置
IotModbusDeviceConfigRespDTO config = configCacheService.getConfig(message.getDeviceId());
if (config == null) {
log.warn("[handle][设备 {} 没有 Modbus 配置]", message.getDeviceId());
return;
}
// 1.3 获取连接信息
ConnectionInfo connInfo = connectionManager.getConnectionInfoByDeviceId(message.getDeviceId());
if (connInfo == null) {
log.warn("[handle][设备 {} 没有连接]", message.getDeviceId());
return;
}
// 2. 解析属性值并写入
Object params = message.getParams();
if (!(params instanceof Map)) {
log.warn("[handle][params 不是 Map 类型: {}]", params);
return;
}
Map