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
package cn.iocoder.yudao.module.iot.gateway.protocol.tcp;
 
import cn.hutool.extra.spring.SpringUtil;
import cn.iocoder.yudao.module.iot.core.enums.IotProtocolTypeEnum;
import cn.iocoder.yudao.module.iot.core.enums.IotSerializeTypeEnum;
import cn.iocoder.yudao.module.iot.core.messagebus.core.IotMessageBus;
import cn.iocoder.yudao.module.iot.core.util.IotDeviceMessageUtils;
import cn.iocoder.yudao.module.iot.gateway.config.IotGatewayProperties;
import cn.iocoder.yudao.module.iot.gateway.config.IotGatewayProperties.ProtocolProperties;
import cn.iocoder.yudao.module.iot.gateway.protocol.IotProtocol;
import cn.iocoder.yudao.module.iot.gateway.protocol.tcp.codec.IotTcpFrameCodec;
import cn.iocoder.yudao.module.iot.gateway.protocol.tcp.codec.IotTcpFrameCodecFactory;
import cn.iocoder.yudao.module.iot.gateway.protocol.tcp.handler.downstream.IotTcpDownstreamHandler;
import cn.iocoder.yudao.module.iot.gateway.protocol.tcp.handler.downstream.IotTcpDownstreamSubscriber;
import cn.iocoder.yudao.module.iot.gateway.protocol.tcp.handler.upstream.IotTcpUpstreamHandler;
import cn.iocoder.yudao.module.iot.gateway.protocol.tcp.manager.IotTcpConnectionManager;
import cn.iocoder.yudao.module.iot.gateway.serialize.IotMessageSerializer;
import cn.iocoder.yudao.module.iot.gateway.serialize.IotMessageSerializerManager;
import io.vertx.core.Vertx;
import io.vertx.core.net.NetServer;
import io.vertx.core.net.NetServerOptions;
import io.vertx.core.net.PemKeyCertOptions;
import lombok.Getter;
import lombok.extern.slf4j.Slf4j;
import cn.hutool.core.lang.Assert;
 
/**
 * IoT TCP 协议实现
 * <p>
 * 基于 Vert.x 实现 TCP 服务器,接收设备上行消息
 *
 * @author 芋道源码
 */
@Slf4j
public class IotTcpProtocol implements IotProtocol {
 
    /**
     * 协议配置
     */
    private final ProtocolProperties properties;
    /**
     * 服务器 ID(用于消息追踪,全局唯一)
     */
    @Getter
    private final String serverId;
 
    /**
     * 运行状态
     */
    @Getter
    private volatile boolean running = false;
 
    /**
     * Vert.x 实例
     */
    private Vertx vertx;
    /**
     * TCP 服务器
     */
    private NetServer tcpServer;
    /**
     * TCP 连接管理器
     */
    private final IotTcpConnectionManager connectionManager;
 
    /**
     * 下行消息订阅者
     */
    private IotTcpDownstreamSubscriber downstreamSubscriber;
 
    /**
     * 消息序列化器
     */
    private final IotMessageSerializer serializer;
    /**
     * TCP 帧编解码器
     */
    private final IotTcpFrameCodec frameCodec;
 
    public IotTcpProtocol(ProtocolProperties properties) {
        IotTcpConfig tcpConfig = properties.getTcp();
        Assert.notNull(tcpConfig, "TCP 协议配置(tcp)不能为空");
        Assert.notNull(tcpConfig.getCodec(), "TCP 拆包配置(tcp.codec)不能为空");
        this.properties = properties;
        this.serverId = IotDeviceMessageUtils.generateServerId(properties.getPort());
 
        // 初始化序列化器
        IotSerializeTypeEnum serializeType = IotSerializeTypeEnum.of(properties.getSerialize());
        Assert.notNull(serializeType, "不支持的序列化类型:" + properties.getSerialize());
        IotMessageSerializerManager serializerManager = SpringUtil.getBean(IotMessageSerializerManager.class);
        this.serializer = serializerManager.get(serializeType);
        // 初始化帧编解码器
        this.frameCodec = IotTcpFrameCodecFactory.create(tcpConfig.getCodec());
 
        // 初始化连接管理器
        this.connectionManager = new IotTcpConnectionManager(tcpConfig.getMaxConnections());
    }
 
    @Override
    public String getId() {
        return properties.getId();
    }
 
    @Override
    public IotProtocolTypeEnum getType() {
        return IotProtocolTypeEnum.TCP;
    }
 
    @Override
    public void start() {
        if (running) {
            log.warn("[start][IoT TCP 协议 {} 已经在运行中]", getId());
            return;
        }
 
        // 1.1 创建 Vertx 实例
        this.vertx = Vertx.vertx();
 
        // 1.2 创建服务器选项
        IotTcpConfig tcpConfig = properties.getTcp();
        NetServerOptions options = new NetServerOptions()
                .setPort(properties.getPort())
                .setTcpKeepAlive(true)
                .setTcpNoDelay(true)
                .setReuseAddress(true)
                .setIdleTimeout((int) (tcpConfig.getKeepAliveTimeoutMs() / 1000)); // 设置空闲超时
        IotGatewayProperties.SslConfig sslConfig = properties.getSsl();
        if (sslConfig != null && Boolean.TRUE.equals(sslConfig.getSsl())) {
            PemKeyCertOptions pemKeyCertOptions = new PemKeyCertOptions()
                    .setKeyPath(sslConfig.getSslKeyPath())
                    .setCertPath(sslConfig.getSslCertPath());
            options.setSsl(true).setKeyCertOptions(pemKeyCertOptions);
        }
 
        // 1.3 创建服务器并设置连接处理器
        tcpServer = vertx.createNetServer(options);
        tcpServer.connectHandler(socket -> {
            IotTcpUpstreamHandler handler = new IotTcpUpstreamHandler(serverId, frameCodec, serializer, connectionManager);
            handler.handle(socket);
        });
 
        // 1.4 启动 TCP 服务器
        try {
            tcpServer.listen().result();
            running = true;
            log.info("[start][IoT TCP 协议 {} 启动成功,端口:{},serverId:{}]",
                    getId(), properties.getPort(), serverId);
 
            // 2. 启动下行消息订阅者
            IotTcpDownstreamHandler downstreamHandler = new IotTcpDownstreamHandler(connectionManager, frameCodec, serializer);
            IotMessageBus messageBus = SpringUtil.getBean(IotMessageBus.class);
            this.downstreamSubscriber = new IotTcpDownstreamSubscriber(this, downstreamHandler, messageBus);
            this.downstreamSubscriber.start();
        } catch (Exception e) {
            log.error("[start][IoT TCP 协议 {} 启动失败]", getId(), e);
            stop0();
            throw e;
        }
    }
 
    @Override
    public void stop() {
        if (!running) {
            return;
        }
        stop0();
    }
 
    private void stop0() {
        // 1. 停止下行消息订阅者
        if (downstreamSubscriber != null) {
            try {
                downstreamSubscriber.stop();
                log.info("[stop][IoT TCP 协议 {} 下行消息订阅者已停止]", getId());
            } catch (Exception e) {
                log.error("[stop][IoT TCP 协议 {} 下行消息订阅者停止失败]", getId(), e);
            }
            downstreamSubscriber = null;
        }
 
        // 2.1 关闭所有连接
        connectionManager.closeAll();
        // 2.2 关闭 TCP 服务器
        if (tcpServer != null) {
            try {
                tcpServer.close().result();
                log.info("[stop][IoT TCP 协议 {} 服务器已停止]", getId());
            } catch (Exception e) {
                log.error("[stop][IoT TCP 协议 {} 服务器停止失败]", getId(), e);
            }
            tcpServer = null;
        }
        // 2.3 关闭 Vertx 实例
        if (vertx != null) {
            try {
                vertx.close().result();
                log.info("[stop][IoT TCP 协议 {} Vertx 已关闭]", getId());
            } catch (Exception e) {
                log.error("[stop][IoT TCP 协议 {} Vertx 关闭失败]", getId(), e);
            }
            vertx = null;
        }
        running = false;
        log.info("[stop][IoT TCP 协议 {} 已停止]", getId());
    }
 
}