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
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
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
package cn.iocoder.yudao.module.iot.gateway.protocol.modbus.tcpserver.codec;
 
import cn.iocoder.yudao.module.iot.core.enums.modbus.IotModbusFrameFormatEnum;
import cn.iocoder.yudao.module.iot.gateway.protocol.modbus.common.utils.IotModbusCommonUtils;
import io.vertx.core.Handler;
import io.vertx.core.buffer.Buffer;
import io.vertx.core.parsetools.RecordParser;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
 
import java.nio.ByteBuffer;
import java.nio.ByteOrder;
import java.nio.charset.StandardCharsets;
import java.util.function.BiConsumer;
 
/**
 * IoT Modbus 帧解码器:集成 TCP 拆包 + 帧格式探测 + 帧解码,一条龙完成从 TCP 字节流到 IotModbusFrame 的转换。
 * <p>
 * 流程:
 * 1. 首帧检测:读前 6 字节,判断 MODBUS_TCP(ProtocolId==0x0000 且 Length 合理)或 MODBUS_RTU
 * 2. 检测后切换到对应的拆包 Handler,并将首包 6 字节通过 handleFirstBytes() 交给新 Handler 处理
 * 3. 拆包完成后解码为 IotModbusFrame,通过回调返回
 *      - MODBUS_TCP:两阶段 RecordParser(MBAP length 字段驱动)
 *      - MODBUS_RTU:功能码驱动的状态机
 *
 * @author 芋道源码
 */
@RequiredArgsConstructor
@Slf4j
public class IotModbusFrameDecoder {
 
    private static final Boolean REQUEST_MODE_DEFAULT = false;
 
    /**
     * 自定义功能码
     */
    private final int customFunctionCode;
 
    /**
     * 创建带自动帧格式检测的 RecordParser(默认响应模式)
     *
     * @param frameHandler 完整帧回调(解码后的 IotModbusFrame + 检测到的帧格式)
     * @return RecordParser 实例
     */
    public RecordParser createRecordParser(BiConsumer<IotModbusFrame, IotModbusFrameFormatEnum> frameHandler) {
        return createRecordParser(frameHandler, REQUEST_MODE_DEFAULT);
    }
 
    /**
     * 创建带自动帧格式检测的 RecordParser
     *
     * @param frameHandler 完整帧回调(解码后的 IotModbusFrame + 检测到的帧格式)
     * @param requestMode  是否为请求模式(true:接收方收到的是 Modbus 请求帧,FC01-04 按固定 8 字节解析;
     *                     false:接收方收到的是 Modbus 响应帧,FC01-04 按 byteCount 变长解析)
     * @return RecordParser 实例
     */
    public RecordParser createRecordParser(BiConsumer<IotModbusFrame, IotModbusFrameFormatEnum> frameHandler,
                                           boolean requestMode) {
        // 先创建一个 RecordParser:使用 fixedSizeMode(6) 读取首帧前 6 字节进行帧格式检测
        RecordParser parser = RecordParser.newFixed(6);
        parser.handler(new DetectPhaseHandler(parser, customFunctionCode, frameHandler, requestMode));
        return parser;
    }
 
    // ==================== 帧解码 ====================
 
    /**
     * 解码响应帧(拆包后的完整帧 byte[])
     *
     * @param data   完整帧字节数组
     * @param format 帧格式
     * @return 解码后的 IotModbusFrame
     */
    private IotModbusFrame decodeResponse(byte[] data, IotModbusFrameFormatEnum format) {
        if (format == IotModbusFrameFormatEnum.MODBUS_TCP) {
            return decodeTcpResponse(data);
        } else {
            return decodeRtuResponse(data);
        }
    }
 
    /**
     * 解码 MODBUS_TCP 响应
     * 格式:[TransactionId(2)] [ProtocolId(2)] [Length(2)] [UnitId(1)] [FC(1)] [Data...]
     */
    private IotModbusFrame decodeTcpResponse(byte[] data) {
        if (data.length < 8) {
            log.warn("[decodeTcpResponse][数据长度不足: {}]", data.length);
            return null;
        }
        ByteBuffer buf = ByteBuffer.wrap(data).order(ByteOrder.BIG_ENDIAN);
        int transactionId = buf.getShort() & 0xFFFF;
        buf.getShort(); // protocolId:固定 0x0000,Modbus 协议标识
        buf.getShort(); // length:后续字节数(UnitId + PDU),拆包阶段已使用
        int slaveId = buf.get() & 0xFF;
        int functionCode = buf.get() & 0xFF;
        // 提取 PDU 数据(从 functionCode 之后到末尾)
        byte[] pdu = new byte[data.length - 8];
        System.arraycopy(data, 8, pdu, 0, pdu.length);
        // 构建 IotModbusFrame
        return buildFrame(slaveId, functionCode, pdu, transactionId);
    }
 
    /**
     * 解码 MODBUS_RTU 响应
     * 格式:[SlaveId(1)] [FC(1)] [Data...] [CRC(2)]
     */
    private IotModbusFrame decodeRtuResponse(byte[] data) {
        if (data.length < 4) {
            log.warn("[decodeRtuResponse][数据长度不足: {}]", data.length);
            return null;
        }
        // 校验 CRC
        if (!IotModbusCommonUtils.verifyCrc16(data)) {
            log.warn("[decodeRtuResponse][CRC 校验失败]");
            return null;
        }
        int slaveId = data[0] & 0xFF;
        int functionCode = data[1] & 0xFF;
        // PDU 数据(不含 slaveId、functionCode、CRC)
        byte[] pdu = new byte[data.length - 4];
        System.arraycopy(data, 2, pdu, 0, pdu.length);
        // 构建 IotModbusFrame
        return buildFrame(slaveId, functionCode, pdu, null);
    }
 
    /**
     * 构建 IotModbusFrame
     */
    private IotModbusFrame buildFrame(int slaveId, int functionCode, byte[] pdu, Integer transactionId) {
        IotModbusFrame frame = new IotModbusFrame()
                .setSlaveId(slaveId)
                .setFunctionCode(functionCode)
                .setPdu(pdu)
                .setTransactionId(transactionId);
        // 异常响应
        if (IotModbusCommonUtils.isExceptionResponse(functionCode)) {
            frame.setFunctionCode(IotModbusCommonUtils.extractOriginalFunctionCode(functionCode));
            if (pdu.length >= 1) {
                frame.setExceptionCode(pdu[0] & 0xFF);
            }
            return frame;
        }
        // 自定义功能码
        if (functionCode == customFunctionCode) {
            // data 区格式:[byteCount(1)] [JSON data(N)]
            if (pdu.length >= 1) {
                int byteCount = pdu[0] & 0xFF;
                if (pdu.length >= 1 + byteCount) {
                    frame.setCustomData(new String(pdu, 1, byteCount, StandardCharsets.UTF_8));
                }
            }
        }
        return frame;
    }
 
    // ==================== 拆包 Handler ====================
 
    /**
     * 帧格式检测阶段 Handler(仅处理首包,探测后切换到对应的拆包 Handler)
     */
    @RequiredArgsConstructor
    private class DetectPhaseHandler implements Handler<Buffer> {
 
        private final RecordParser parser;
        private final int customFunctionCode;
        private final BiConsumer<IotModbusFrame, IotModbusFrameFormatEnum> frameHandler;
        private final boolean requestMode;
 
        @Override
        public void handle(Buffer buffer) {
            // 检测帧格式:protocolId==0x0000 且 length 合法 → MODBUS_TCP,否则 → MODBUS_RTU
            byte[] bytes = buffer.getBytes();
            int protocolId = ((bytes[2] & 0xFF) << 8) | (bytes[3] & 0xFF);
            int length = ((bytes[4] & 0xFF) << 8) | (bytes[5] & 0xFF);
 
            // 分别处理 MODBUS_TCP、MODBUS_RTU 两种情况
            if (protocolId == 0x0000 && length >= 1 && length <= 253) {
                // MODBUS_TCP:切换到 TCP 拆包 Handler
                log.debug("[DetectPhaseHandler][检测到 MODBUS_TCP 帧格式]");
                TcpFrameHandler tcpHandler = new TcpFrameHandler(parser, frameHandler);
                parser.handler(tcpHandler);
                // 当前 bytes 就是 MBAP 的前 6 字节,直接交给 tcpHandler 处理
                tcpHandler.handleFirstBytes(bytes);
            } else {
                // MODBUS_RTU:切换到 RTU 拆包 Handler
                log.debug("[DetectPhaseHandler][检测到 MODBUS_RTU 帧格式]");
                RtuFrameHandler rtuHandler = new RtuFrameHandler(parser, frameHandler, customFunctionCode, requestMode);
                parser.handler(rtuHandler);
                // 当前 bytes 包含前 6 字节(slaveId + FC + 部分数据),交给 rtuHandler 处理
                rtuHandler.handleFirstBytes(bytes);
            }
        }
    }
 
    /**
     * MODBUS_TCP 拆包 Handler(两阶段 RecordParser)
     * <p>
     * Phase 1: fixedSizeMode(6) → 读 MBAP 前 6 字节,提取 length
     * Phase 2: fixedSizeMode(length) → 读 unitId + PDU
     */
    @RequiredArgsConstructor
    private class TcpFrameHandler implements Handler<Buffer> {
 
        private final RecordParser parser;
        private final BiConsumer<IotModbusFrame, IotModbusFrameFormatEnum> frameHandler;
 
        private byte[] mbapHeader;
        private boolean waitingForBody = false;
 
        /**
         * 处理探测阶段传来的首帧 6 字节(即 MBAP 头)
         *
         * @param bytes 探测阶段消费的 6 字节
         */
        void handleFirstBytes(byte[] bytes) {
            int length = ((bytes[4] & 0xFF) << 8) | (bytes[5] & 0xFF);
            this.mbapHeader = bytes;
            this.waitingForBody = true;
            parser.fixedSizeMode(length);
        }
 
        @Override
        public void handle(Buffer buffer) {
            if (waitingForBody) {
                // Phase 2: 收到 body(unitId + PDU)
                byte[] body = buffer.getBytes();
                // 拼接完整帧:MBAP(6) + body
                byte[] fullFrame = new byte[mbapHeader.length + body.length];
                System.arraycopy(mbapHeader, 0, fullFrame, 0, mbapHeader.length);
                System.arraycopy(body, 0, fullFrame, mbapHeader.length, body.length);
                // 解码并回调
                IotModbusFrame frame = decodeResponse(fullFrame, IotModbusFrameFormatEnum.MODBUS_TCP);
                if (frame != null) {
                    frameHandler.accept(frame, IotModbusFrameFormatEnum.MODBUS_TCP);
                }
                // 切回 Phase 1
                waitingForBody = false;
                mbapHeader = null;
                parser.fixedSizeMode(6);
            } else {
                // Phase 1: 收到 MBAP 头 6 字节
                byte[] header = buffer.getBytes();
                int length = ((header[4] & 0xFF) << 8) | (header[5] & 0xFF);
                if (length < 1 || length > 253) {
                    log.warn("[TcpFrameHandler][MBAP Length 异常: {}]", length);
                    parser.fixedSizeMode(6);
                    return;
                }
                this.mbapHeader = header;
                this.waitingForBody = true;
                parser.fixedSizeMode(length);
            }
        }
    }
 
    /**
     * MODBUS_RTU 拆包 Handler(功能码驱动的状态机)
     * <p>
     * 状态机流程:
     * Phase 1: fixedSizeMode(2) → 读 slaveId + functionCode
     * Phase 2: 根据 functionCode 确定剩余长度:
     * - 异常响应 (FC & EXCEPTION_MASK):fixedSizeMode(3) → exceptionCode(1) + CRC(2)
     * - 自定义 FC / FC01-04 响应:fixedSizeMode(1) → 读 byteCount → fixedSizeMode(byteCount + 2)
     * - FC05/06 响应:fixedSizeMode(6) → addr(2) + value(2) + CRC(2)
     * - FC15/16 响应:fixedSizeMode(6) → addr(2) + quantity(2) + CRC(2)
     * <p>
     * 请求模式(requestMode=true)时,FC01-04 按固定 8 字节解析(与写响应相同路径),
     * 因为读请求格式为 [SlaveId(1)][FC(1)][StartAddr(2)][Quantity(2)][CRC(2)]
     */
    @RequiredArgsConstructor
    private class RtuFrameHandler implements Handler<Buffer> {
 
        private static final int STATE_HEADER = 0;
        private static final int STATE_EXCEPTION_BODY = 1;
        private static final int STATE_READ_BYTE_COUNT = 2;
        private static final int STATE_READ_DATA = 3;
        private static final int STATE_WRITE_BODY = 4;
 
        private final RecordParser parser;
        private final BiConsumer<IotModbusFrame, IotModbusFrameFormatEnum> frameHandler;
        private final int customFunctionCode;
        /**
         * 请求模式:
         *    - true 表示接收方收到的是 Modbus 请求帧(如设备端收到网关下发的读请求),FC01-04 按固定 8 字节帧解析
         *    - false 表示接收方收到的是 Modbus 响应帧,FC01-04 按 byteCount 变长解析
         */
        private final boolean requestMode;
 
        private int state = STATE_HEADER;
        private byte slaveId;
        private byte functionCode;
        private byte byteCount;
        private Buffer pendingData;
        private int expectedDataLen;
 
        /**
         * 处理探测阶段传来的首帧 6 字节
         * <p>
         * 由于 RTU 首帧被探测阶段消费了 6 字节,这里需要从中提取 slaveId + FC 并根据 FC 处理剩余数据
         *
         * @param bytes 探测阶段消费的 6 字节:[slaveId][FC][...4 bytes...]
         */
        void handleFirstBytes(byte[] bytes) {
            this.slaveId = bytes[0];
            this.functionCode = bytes[1];
            int fc = functionCode & 0xFF;
            if (IotModbusCommonUtils.isExceptionResponse(fc)) {
                // 异常响应:完整帧 = slaveId(1) + FC(1) + exceptionCode(1) + CRC(2) = 5 字节
                // 已有 6 字节(多 1 字节),取前 5 字节组装
                Buffer frame = Buffer.buffer(5);
                frame.appendByte(slaveId);
                frame.appendByte(functionCode);
                frame.appendBytes(bytes, 2, 3); // exceptionCode + CRC
                emitFrame(frame);
                resetToHeader();
            } else if (IotModbusCommonUtils.isReadResponse(fc) && requestMode) {
                // 请求模式下的读请求:固定 8 字节 [SlaveId(1)][FC(1)][StartAddr(2)][Quantity(2)][CRC(2)]
                // 已有 6 字节,还需 2 字节(CRC)
                state = STATE_WRITE_BODY;
                this.pendingData = Buffer.buffer();
                this.pendingData.appendBytes(bytes, 2, 4); // 暂存已有的 4 字节(StartAddr + Quantity)
                parser.fixedSizeMode(2); // 还需 2 字节(CRC)
            } else if (IotModbusCommonUtils.isReadResponse(fc) || fc == customFunctionCode) {
                // 读响应或自定义 FC:bytes[2] = byteCount
                this.byteCount = bytes[2];
                int bc = byteCount & 0xFF;
                // 已有数据:bytes[3..5] = 3 字节
                // 还需:byteCount + CRC(2) - 3 字节已有
                int remaining = bc + 2 - 3;
                if (remaining <= 0) {
                    // 数据已足够,组装完整帧
                    int totalLen = 2 + 1 + bc + 2; // slaveId + FC + byteCount + data + CRC
                    Buffer frame = Buffer.buffer(totalLen);
                    frame.appendByte(slaveId);
                    frame.appendByte(functionCode);
                    frame.appendByte(byteCount);
                    frame.appendBytes(bytes, 3, bc + 2); // data + CRC
                    emitFrame(frame);
                    resetToHeader();
                } else {
                    // 需要继续读
                    state = STATE_READ_DATA;
                    this.pendingData = Buffer.buffer();
                    this.pendingData.appendBytes(bytes, 3, 3); // 暂存已有的 3 字节
                    this.expectedDataLen = bc + 2; // byteCount 个数据 + 2 CRC
                    parser.fixedSizeMode(remaining);
                }
            } else if (IotModbusCommonUtils.isWriteResponse(fc)) {
                // 写响应:总长 = slaveId(1) + FC(1) + addr(2) + value/qty(2) + CRC(2) = 8 字节
                // 已有 6 字节,还需 2 字节
                state = STATE_WRITE_BODY;
                this.pendingData = Buffer.buffer();
                this.pendingData.appendBytes(bytes, 2, 4); // 暂存已有的 4 字节
                parser.fixedSizeMode(2); // 还需 2 字节(CRC)
            } else {
                log.warn("[RtuFrameHandler][未知功能码: 0x{}]", Integer.toHexString(fc));
                resetToHeader();
            }
        }
 
        @Override
        public void handle(Buffer buffer) {
            switch (state) {
                case STATE_HEADER:
                    handleHeader(buffer);
                    break;
                case STATE_EXCEPTION_BODY:
                    handleExceptionBody(buffer);
                    break;
                case STATE_READ_BYTE_COUNT:
                    handleReadByteCount(buffer);
                    break;
                case STATE_READ_DATA:
                    handleReadData(buffer);
                    break;
                case STATE_WRITE_BODY:
                    handleWriteBody(buffer);
                    break;
                default:
                    resetToHeader();
            }
        }
 
        private void handleHeader(Buffer buffer) {
            byte[] header = buffer.getBytes();
            this.slaveId = header[0];
            this.functionCode = header[1];
            int fc = functionCode & 0xFF;
            if (IotModbusCommonUtils.isExceptionResponse(fc)) {
                // 异常响应
                state = STATE_EXCEPTION_BODY;
                parser.fixedSizeMode(3); // exceptionCode(1) + CRC(2)
            } else if (IotModbusCommonUtils.isReadResponse(fc) && requestMode) {
                // 请求模式下的读请求:固定 8 字节,已读 2 字节(slaveId + FC),还需 6 字节
                state = STATE_WRITE_BODY;
                pendingData = Buffer.buffer();
                parser.fixedSizeMode(6); // StartAddr(2) + Quantity(2) + CRC(2)
            } else if (IotModbusCommonUtils.isReadResponse(fc) || fc == customFunctionCode) {
                // 读响应或自定义 FC
                state = STATE_READ_BYTE_COUNT;
                parser.fixedSizeMode(1); // byteCount
            } else if (IotModbusCommonUtils.isWriteResponse(fc)) {
                // 写响应
                state = STATE_WRITE_BODY;
                pendingData = Buffer.buffer();
                parser.fixedSizeMode(6); // addr(2) + value(2) + CRC(2)
            } else {
                log.warn("[RtuFrameHandler][未知功能码: 0x{}]", Integer.toHexString(fc));
                resetToHeader();
            }
        }
 
        private void handleExceptionBody(Buffer buffer) {
            // buffer = exceptionCode(1) + CRC(2)
            Buffer frame = Buffer.buffer();
            frame.appendByte(slaveId);
            frame.appendByte(functionCode);
            frame.appendBuffer(buffer);
            emitFrame(frame);
            resetToHeader();
        }
 
        private void handleReadByteCount(Buffer buffer) {
            this.byteCount = buffer.getByte(0);
            int bc = byteCount & 0xFF;
            state = STATE_READ_DATA;
            pendingData = Buffer.buffer();
            expectedDataLen = bc + 2; // data(bc) + CRC(2)
            parser.fixedSizeMode(expectedDataLen);
        }
 
        private void handleReadData(Buffer buffer) {
            pendingData.appendBuffer(buffer);
            if (pendingData.length() >= expectedDataLen) {
                // 组装完整帧
                Buffer frame = Buffer.buffer();
                frame.appendByte(slaveId);
                frame.appendByte(functionCode);
                frame.appendByte(byteCount);
                frame.appendBuffer(pendingData);
                emitFrame(frame);
                resetToHeader();
            }
            // 否则继续等待(不应该发生,因为我们精确设置了 fixedSizeMode)
        }
 
        private void handleWriteBody(Buffer buffer) {
            pendingData.appendBuffer(buffer);
            // 完整帧
            Buffer frame = Buffer.buffer();
            frame.appendByte(slaveId);
            frame.appendByte(functionCode);
            frame.appendBuffer(pendingData);
            emitFrame(frame);
            resetToHeader();
        }
 
        /**
         * 发射完整帧:解码并回调
         */
        private void emitFrame(Buffer frameBuffer) {
            IotModbusFrame frame = decodeResponse(frameBuffer.getBytes(), IotModbusFrameFormatEnum.MODBUS_RTU);
            if (frame != null) {
                frameHandler.accept(frame, IotModbusFrameFormatEnum.MODBUS_RTU);
            }
        }
 
        private void resetToHeader() {
            state = STATE_HEADER;
            pendingData = null;
            parser.fixedSizeMode(2); // slaveId + FC
        }
 
    }
 
}