ControlRookie
返回文章

第2篇_CONNECT 和 CONNACK 怎么读?我带你把十六进制和 ST 代码对上

这一篇正式进入 MQTT 握手核心,重点讲清 CONNECT 和 CONNACK 的报文结构、十六进制拆解、状态机推进关系,以及它们在 MqttClient_V1_0 里的真实 ST 落点。

这一篇正式进入 MQTT 握手核心,重点讲清 CONNECT 和 CONNACK 的报文结构、十六进制拆解、状态机推进关系,以及它们在 MqttClient_V1_0 里的真实 ST 落点。

适合谁收藏

  • 正在做 CODESYS / PLC / MQTT 项目的人
  • 想把 MQTT 从报文真正看到 ST 代码的人
  • 正在排查 QoS1 / QoS2 超时、掉线、重连问题的人

上一篇我们只干了一件事: 把 MQTT 在协议栈里的位置钉死了。

这一篇开始,正式进核心。

这一篇只解决 3 个问题:

  1. MQTT 连接为什么不是 TCP 通了就结束
  2. CONNECT / CONNACK 报文到底长什么样
  3. 在 MqttClient_V1_0 里,这两类报文是怎么落到 ST 代码里的

先给结论:

TCP 连上,只说明“路打通了”。

一、先看完整握手链

先别急着看字节。 先看一次完整握手到底发生了什么。

Mermaid
sequenceDiagram
    participant App as 用户逻辑
    participant PLC as FB_MqttClient
    participant TCP as NBS.TCP_Client
    participant Broker as MQTT Broker

    App->>PLC: bEnable=TRUE, bConnect 上升沿
    PLC->>PLC: iDisconnected -> iTcpConnect
    PLC->>TCP: M_TcpClient()
    TCP-->>PLC: TCP connected
    PLC->>PLC: iTcpConnect -> iConnect
    PLC->>PLC: M_BuildConnectPacket()
    PLC->>Broker: CONNECT
    PLC->>PLC: xWaitingForAck=TRUE
    PLC->>PLC: byExpectedMsgType=CONNACK
    PLC->>PLC: iConnect -> iConnAck
    Broker-->>PLC: CONNACK
    PLC->>PLC: M_HandleConnAck()
    PLC->>PLC: 保存服务端能力
    PLC->>PLC: iConnAck -> iConnected

这张图先记住两句话:

  1. iTcpConnect 只是 TCP 层,iConnect 才真正进入 MQTT 层。
  2. iConnAck 不是“等个回包”这么简单,而是在等 Broker 对本次连接的正式裁决。

二、CONNECT 报文到底在干嘛

CONNECT 报文本质上是在告诉 Broker:

  • 我是谁
  • 我要用哪个 MQTT 版本
  • 我的 KeepAlive 是多少
  • 我要不要 Clean Session / Clean Start
  • 我有没有用户名密码
  • 我有没有遗嘱消息
  • 如果是 MQTT 5.0,我还有哪些连接属性

也就是说,CONNECT 不是“你好我来了”这么简单。 它其实是一份比较完整的“连接申请表”。


三、CONNECT 报文结构先看表

1. 固定报头

字段说明
报文类型0x10,表示 CONNECT
Remaining Length后面所有字节的总长度

2. 可变报头

字段说明
Protocol Name一般是 MQTT
Protocol Level3.1.1 对应 0x04,5.0 对应 0x05
Connect Flags用户名、密码、遗嘱、Clean Session 等都在这里
Keep Alive心跳周期,单位秒
Properties仅 MQTT 5.0 有

3. 载荷

字段说明
Client ID客户端标识
Will Properties仅 MQTT 5.0 且 Will Flag=1 时存在
Will Topic遗嘱主题
Will Payload遗嘱内容
Username可选
Password可选

四、先拆一个最典型的 CONNECT 十六进制

我们先看一个最小可用版 MQTT 3.1.1 CONNECT:

text
10 13 00 04 4D 51 54 54 04 02 00 3C 00 07 70 6C 63 5F 30 30 31

拆开看:

字节含义说明
10CONNECT高四位是报文类型
13Remaining Length后面共有 19 字节
00 04协议名长度MQTT 长度 4
4D 51 54 54M Q T T协议名
04协议版本MQTT 3.1.1
02Connect FlagsClean Session=1
00 3CKeepAlive60 秒
00 07Client ID 长度7
70 6C 63 5F 30 30 31plc_001Client ID

这里最关键的是第 10 个字节,也就是 Connect Flags。


五、Connect Flags 为什么很关键

Connect Flags 是个单字节,但是信息密度很高。

Bit含义
bit7Username Flag
bit6Password Flag
bit5Will Retain
bit4-bit3Will QoS
bit2Will Flag
bit1Clean Session
bit0保留位,必须为 0

如果这 1 个字节拼错,Broker 直接就能拒绝你。

这也是为什么源码里没有偷懒,而是在构包前就先做参数一致性校验。

看真实 ST 代码:

iecst
IF (sPassword <> '') AND (sUsername = '') THEN
    M_SetError(
        uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
        sMessage := 'Password requires username');
    M_BuildConnectPacket := FALSE;
    RETURN;
END_IF

IF (NOT bWillFlag) AND ((sWillTopic <> '') OR (sWillMessage <> '')) THEN
    M_SetError(
        uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
        sMessage := 'Will fields require Will Flag');
    M_BuildConnectPacket := FALSE;
    RETURN;
END_IF

这段代码很值钱。 它不是为了“写得严谨一点”,而是为了避免你发出去的 CONNECT 先天就不合法。


六、M_BuildConnectPacket 到底怎么拼包

这段方法的主线其实很清楚,可以压成 4 步:

Mermaid
flowchart TD
    A[参数一致性校验] --> B[计算可变报头和载荷长度]
    B --> C[写固定报头和 Remaining Length]
    C --> D[写协议名 版本 Flags KeepAlive]
    D --> E[如果是 MQTT 5.0 再写连接属性]
    E --> F[写 Client ID Will Username Password]

源码里最核心的一段就是协议头拼接:

iecst
aTxBuf[0] := E_MqttPacketType.byConnect;
uiPos := uiPos + 1;

uiPos := uiPos + M_EncodeRemainingLength(uiRemainingLen, ADR(aTxBuf[uiPos]));

aTxBuf[uiPos] := 0;
uiPos := uiPos + 1;
aTxBuf[uiPos] := 4;
uiPos := uiPos + 1;
aTxBuf[uiPos] := 77;
uiPos := uiPos + 1;
aTxBuf[uiPos] := 81;
uiPos := uiPos + 1;
aTxBuf[uiPos] := 84;
uiPos := uiPos + 1;
aTxBuf[uiPos] := 84;
uiPos := uiPos + 1;

你不用纠结 77 81 84 84 这种写法看起来有点笨。 它的好处恰恰是:

  • 明确
  • 不依赖额外库
  • 在 PLC 里非常可控

七、MQTT 5.0 比 3.1.1 多了什么

如果是 MQTT 5.0,CONNECT 里会多一块 Properties。

在这个开源库里,当前会构造的连接属性包括:

属性作用
Session Expiry Interval会话过期时间
Receive Maximum声明自己接收窗口
Maximum Packet Size声明自己可接收的最大报文
Topic Alias Maximum声明自己支持的主题别名上限
Request Problem Information请求问题信息
Request Response Information请求响应信息

源码里能直接看到这段长度计算:

iecst
IF eVersion = E_MqttVersion.byMqttVersion50 THEN
    IF udiSessionExpiry > 0 THEN
        uiPropsLen := uiPropsLen + 5;
    END_IF
    uiPropsLen := uiPropsLen + 3;
    uiPropsLen := uiPropsLen + 5;
    uiPropsLen := uiPropsLen + 3;
    uiPropsLen := uiPropsLen + 2;
    IF bRequestResponseInfo THEN
        uiPropsLen := uiPropsLen + 2;
    END_IF
END_IF

这段背后的工程意义是:

MQTT 5.0 不是多几个花哨字段。

八、Broker 回的 CONNACK 不是“收到就算完”

客户端发完 CONNECT 后,状态机不会直接宣布成功。 它会进入 iConnAck,并且明确设置:

iecst
xWaitingForAck := TRUE
byExpectedMsgType := CONNACK

也就是说,后面收上来的第一帧 MQTT 报文,必须满足:

  1. 类型对
  2. 格式对
  3. 原因码成功
  4. 如果是 5.0,属性也要解析对

否则就不能进 iConnected。


九、CONNACK 报文先看结构

MQTT 3.1.1

字段说明
Acknowledge Flagsbit0 是 Session Present
Return Code是否接受连接

MQTT 5.0

字段说明
Acknowledge Flags会话相关标志
Reason Code成功或失败原因
Properties服务端能力说明

十、先拆一个 CONNACK 十六进制

最小 MQTT 3.1.1 成功 CONNACK:

text
20 02 00 00

拆开看:

字节含义
20CONNACK
02Remaining Length=2
00Acknowledge Flags,Session Present=0
00Return Code=0,连接成功

如果最后一个字节不是 00,那就不是“连上了但状态机没处理好”,而是 Broker 明确拒绝了你。


十一、M_HandleConnAck 到底做了什么

这个方法的职责,不是简单清个标志位。 它至少做了 5 件事:

  1. 检查当前是不是确实在等待 CONNACK
  2. 检查 CONNACK 固定格式是否合法
  3. 解析原因码,失败就输出诊断信息
  4. 如果是 MQTT 5.0,解析服务端属性
  5. 把服务端能力落地保存,供后续发布/订阅使用

先看最前面的合法性判断:

iecst
IF uiRxLength < 4 THEN
    M_HandleConnAck := FALSE;
    RETURN;
END_IF

IF NOT xWaitingForAck OR aRxBuf[0] <> byExpectedMsgType THEN
    M_HandleConnAck := FALSE;
    RETURN;
END_IF

byConnAckFlags := aRxBuf[2];

IF (byConnAckFlags AND 16#FE) <> 0 THEN
    sDiagMsg := 'Invalid CONNACK flags';
    M_HandleConnAck := FALSE;
    RETURN;
END_IF

这段很典型: 不是“只要收到了包就往下跑”,而是一步步卡住协议条件。


十二、MQTT 5.0 真正的价值,在 CONNACK 里才开始体现

M_HandleConnAck 里最值钱的一段,不是判断成功失败,而是解析服务端能力:

变量含义
uiServerReceiveMaxBroker 接收窗口
byServerMaxQoS服务端允许的最大 QoS
bServerRetainAvailable是否支持 Retain
udServerMaxPacketSize最大报文限制
uiServerTopicAliasMax主题别名上限
bServerWildcardSubAvail是否支持通配订阅
bServerSubIdAvail是否支持 Subscription Identifier
bServerSharedSubAvail是否支持共享订阅

也就是说,后面你能不能发某些 MQTT 5.0 特性,不是客户端想发就发。 你得先看 Broker 在 CONNACK 里有没有给你这个权限。


十三、为什么很多“连接成功后的故障”,根其实都在 CONNACK

现场很容易出现一种错觉:

TCP 通了,CONNECT 发了,CONNACK 也回了,那后面肯定没问题。

这句话只对一半。

更准确的说法是:

CONNACK 成功,只说明 Broker 接受了这次连接。

比如:

  • 你想发 QoS2,但 Broker 最大只给到 QoS1
  • 你想发 Retain,但 Broker 不支持
  • 你想用共享订阅或订阅标识符,但 Broker 不开放

这些坑,根本不是后面订阅/发布阶段才突然冒出来的。 它们在 CONNACK 里就已经埋下了。


十四、把标准、报文、状态机、代码串起来看一遍

1. 标准层

  • 客户端发 CONNECT
  • 服务端回 CONNACK

2. 报文层

  • CONNECT 带版本、Flags、KeepAlive、身份信息、5.0 属性
  • CONNACK 带会话标志、原因码、5.0 能力属性

3. 状态机层

text
iDisconnected -> iTcpConnect -> iConnect -> iConnAck -> iConnected

4. ST 实现层

动作方法
构包M_BuildConnectPacket
等待连接确认iConnAck
解析连接确认M_HandleConnAck

这才是这套系列后面一直要坚持的讲法。 不是把它拆成四门课,而是把四层一起看。


十五、这一篇你最该记住的 4 句话

  1. TCP 连上,不等于 MQTT 连上。
  2. CONNECT 是连接申请,CONNACK 是服务端正式裁决。
  3. MQTT 5.0 的很多能力边界,不是在客户端写死的,而是在 CONNACK 里协商出来的。
  4. 真正靠谱的客户端,必须把报文校验、状态机推进和能力保存三件事一起做好。

十六、下篇预告

下一篇我们开始拆最核心的一类报文:

PUBLISH

也就是:

  • 消息到底怎么发出去
  • QoS0 / QoS1 / QoS2 为什么差这么大
  • 报文里的 DUP、QoS、Retain 这几个位到底怎么落到代码里

下一篇开始,代码味会更重。 因为从 PUBLISH 开始,MQTT 才真正进入“能跑”和“跑稳”之间的分水岭。


完整 ST 代码

复制使用说明

  • 这部分给出的是与本篇主题直接对应的完整 ST 代码,不是零碎片段。
  • 如果你只是想先跑通,优先整段复制,不要只摘几行变量或几条赋值语句。
  • 如果是 METHOD,请确认它仍然属于 FB_MqttClient;如果是 PROGRAM,请确认相关 DUT、GVL、FB 已一并导入。

代码阅读重点

  • 先按 报文结构 -> 状态机入口 -> 关键变量 -> 返回结果 的顺序看。
  • 再把正文里的十六进制拆解和这里的字节写入、字节解析语句一行行对上。
  • 最后回到在线调试,重点盯 uiTxLength、uiRxLength、eState、xWaitingForAck 这类状态量。

完整代码 1:M_BuildConnectPacket

  • 对应源码路径:10 MQTT/MqttClient_V1_0/Device/Application/MQTT/POUs/MqttClient NBS/FB_MqttClient/处理发送报文/M_BuildConnectPacket.st
  • 复制使用说明:直接复制到 FB_MqttClient 的同名 METHOD 中即可,对应的是客户端发起连接时的组包逻辑。
  • 阅读重点:重点看固定报头、连接标志字节、KeepAlive、MQTT 5.0 属性长度这几段,它们和 CONNECT 十六进制拆解是一一对应的。
iecst
/// =======================================================================
/// 名称      : M_BuildConnectPacket
/// 功能      : 构建 CONNECT 发送报文
/// 说明      : 按 MQTT 3.1.1/5.0 规则组装 CONNECT 报文并写入发送缓冲区。
/// 编程人员  : ControlRookie
/// 时间      : 2026-05-05
/// 版本      : V1.0
/// =======================================================================
{attribute 'hide_all_locals'}
METHOD M_BuildConnectPacket : BOOL
VAR
    uiPos               : UINT := 0;                                        //字节索引
    uiVarHeaderLen      : UINT;                                             //可变报头长度
    uiPayloadLen        : UINT;                                             //载荷长度
    uiRemainingLen      : UINT;                                             //剩余长度
    uiPropsLen          : UINT;                                             //MQTT 5.0属性长度
    i                   : DINT;
    byConnectFlags      : BYTE;                                             //连接标志字节
    uiTopicAliasMax     : UINT;                                             //客户端接收方向主题别名上限
END_VAR

// === IMPLEMENTATION ===

// BUG-10: 缓冲区溢出保护
IF SIZEOF(aTxBuf) < 256 THEN
    M_BuildConnectPacket := FALSE;
    RETURN;
END_IF

IF (sPassword <> '') AND (sUsername = '') THEN
    M_SetError(
        uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
        sMessage := 'Password requires username');
    M_BuildConnectPacket := FALSE;
    RETURN;
END_IF

IF (NOT bWillFlag) AND ((sWillTopic <> '') OR (sWillMessage <> '')) THEN
    M_SetError(
        uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
        sMessage := 'Will fields require Will Flag');
    M_BuildConnectPacket := FALSE;
    RETURN;
END_IF

IF (bWillRetain OR (eWillQoS <> E_MqttQoS.byQoS0)) AND (NOT bWillFlag) THEN
    M_SetError(
        uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
        sMessage := 'Invalid will flag combination');
    M_BuildConnectPacket := FALSE;
    RETURN;
END_IF

IF bWillFlag AND (sWillTopic = '') THEN
    M_SetError(
        uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
        sMessage := 'Will topic is required');
    M_BuildConnectPacket := FALSE;
    RETURN;
END_IF

// 清空发送缓冲区,避免数据干扰
FOR i := LOWER_BOUND(aTxBuf, 1) TO UPPER_BOUND(aTxBuf, 1) DO
    aTxBuf[i] := 0;
END_FOR

/// =======================================================================
/// 长度计算
/// =======================================================================
// 可变报头长度: 协议名(6) + 协议版本(1) + 连接标志(1) + 保持连接(2)
uiVarHeaderLen := 10;

// MQTT 5.0: 连接属性长度
uiPropsLen := 0;
IF eVersion = E_MqttVersion.byMqttVersion50 THEN
    // Session Expiry Interval (4字节值 + 1字节标识符)
    IF udiSessionExpiry > 0 THEN
        uiPropsLen := uiPropsLen + 5;
    END_IF
    // Receive Maximum (2字节值 + 1字节标识符)
    uiPropsLen := uiPropsLen + 3;
    // Maximum Packet Size (4字节值 + 1字节标识符)
    uiPropsLen := uiPropsLen + 5;
    // Topic Alias Maximum (2字节值 + 1字节标识符)
    uiPropsLen := uiPropsLen + 3;
    // Request Problem Information (1字节值 + 1字节标识符)
    uiPropsLen := uiPropsLen + 2;
    // Request Response Information (1字节值 + 1字节标识符)
    IF bRequestResponseInfo THEN
        uiPropsLen := uiPropsLen + 2;
    END_IF
    // 属性长度字段的编码长度
    IF uiPropsLen < 128 THEN
        uiVarHeaderLen := uiVarHeaderLen + 1 + uiPropsLen;                  //1字节属性长度
    ELSE
        uiVarHeaderLen := uiVarHeaderLen + 2 + uiPropsLen;                  //2字节属性长度
    END_IF
END_IF

// 载荷长度
uiPayloadLen := TO_UINT(2 + LEN(sClientID));                                    //Client ID
IF sUsername <> '' THEN
    uiPayloadLen := uiPayloadLen + TO_UINT(2 + LEN(sUsername));
END_IF
IF sPassword <> '' THEN
    uiPayloadLen := uiPayloadLen + TO_UINT(2 + LEN(sPassword));
END_IF
IF bWillFlag THEN
    IF eVersion = E_MqttVersion.byMqttVersion50 THEN
        uiPayloadLen := uiPayloadLen + 1;
    END_IF
    uiPayloadLen := uiPayloadLen + TO_UINT(2 + LEN(sWillTopic) + 2 + LEN(sWillMessage));
END_IF

IF uiReceiveMax = 0 THEN
    M_SetError(
        uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
        sMessage := 'Receive Maximum must not be zero');
    M_BuildConnectPacket := FALSE;
    RETURN;
END_IF

IF udMaxPacketSize = 0 THEN
    M_SetError(
        uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
        sMessage := 'Maximum Packet Size must not be zero');
    M_BuildConnectPacket := FALSE;
    RETURN;
END_IF

// 剩余长度
uiRemainingLen := uiVarHeaderLen + uiPayloadLen;

// 缓冲区长度检查
IF uiRemainingLen + 5 > SIZEOF(aTxBuf) THEN
    M_BuildConnectPacket := FALSE;
    RETURN;
END_IF

/// =======================================================================
/// 创建报文
/// =======================================================================
uiPos := 0;

// ****************** 固定报文头 = 报文类型 + 剩余长度 ******************
// 报文类型
aTxBuf[0] := E_MqttPacketType.byConnect;
uiPos := uiPos + 1;

// 剩余长度
uiPos := uiPos + M_EncodeRemainingLength(uiRemainingLen, ADR(aTxBuf[uiPos]));

// ****************** 可变报文头 = 协议名(6) + 协议版本(1) + 连接标志(1) + 保持连接(2) ******************
// 协议名:"MQTT"
aTxBuf[uiPos] := 0;                                                         //长度高字节MSB
uiPos := uiPos + 1;
aTxBuf[uiPos] := 4;                                                         //长度低字节LSB
uiPos := uiPos + 1;
aTxBuf[uiPos] := 77;                                                            // 'M'
uiPos := uiPos + 1;
aTxBuf[uiPos] := 81;                                                            // 'Q'
uiPos := uiPos + 1;
aTxBuf[uiPos] := 84;                                                            // 'T'
uiPos := uiPos + 1;
aTxBuf[uiPos] := 84;                                                            // 'T'
uiPos := uiPos + 1;

// 协议版本(根据eVersion选择)
IF eVersion = E_MqttVersion.byMqttVersion50 THEN
    aTxBuf[uiPos] := 5;                                                     //MQTT 5.0
ELSE
    aTxBuf[uiPos] := 4;                                                     //MQTT 3.1.1
END_IF
uiPos := uiPos + 1;

// 连接标志
byConnectFlags :=
        (BOOL_TO_BYTE(FALSE) * 1) OR                                            //Bit 0: Reserved。服务端必须验证保留标志位是否为0
        (BOOL_TO_BYTE(bCleanSession) * 2) OR                                    //Bit 1: Clean Session
        (BOOL_TO_BYTE(bWillFlag AND sWillTopic <> '') * 4) OR                   //Bit 2: Will Flag
        (BOOL_TO_BYTE(eWillQoS = 1 OR eWillQoS = 3) * 8) OR                 //Bit 3: Will QoS低位
        (BOOL_TO_BYTE(eWillQoS = 2 OR eWillQoS = 3) * 16) OR                    //Bit 4: Will QoS高位
        (BOOL_TO_BYTE(bWillFlag AND bWillRetain) * 32) OR                       //Bit 5: Will Retain
        (BOOL_TO_BYTE(sPassword <> '') * 64) OR                             //Bit 6: Password Flag
        (BOOL_TO_BYTE(sUsername <> '') * 128);                                  //Bit 7: Username Flag
aTxBuf[uiPos] := byConnectFlags;
uiPos := uiPos + 1;

// 保持连接
aTxBuf[uiPos] := UINT_TO_BYTE(SHR(uiKeepAlive, 8));                         //KeepAlive高字节
uiPos := uiPos + 1;
aTxBuf[uiPos] := UINT_TO_BYTE(uiKeepAlive AND 16#FF);                           //KeepAlive低字节
uiPos := uiPos + 1;

// ****************** MQTT 5.0: 连接属性 ******************
IF eVersion = E_MqttVersion.byMqttVersion50 THEN
    // 属性长度(可变字节编码)
    IF uiPropsLen < 128 THEN
        aTxBuf[uiPos] := TO_BYTE(uiPropsLen);
        uiPos := uiPos + 1;
    ELSE
        aTxBuf[uiPos] := TO_BYTE((uiPropsLen / 128) OR 16#80);
        uiPos := uiPos + 1;
        aTxBuf[uiPos] := TO_BYTE(uiPropsLen MOD 128);
        uiPos := uiPos + 1;
    END_IF

    // Session Expiry Interval
    IF udiSessionExpiry > 0 THEN
        aTxBuf[uiPos] := GVL_Mqtt.cnPropSessionExpiry;  uiPos := uiPos + 1;
        aTxBuf[uiPos] := UDINT_TO_BYTE(SHR(udiSessionExpiry, 24) AND 16#FF);    uiPos := uiPos + 1;
        aTxBuf[uiPos] := UDINT_TO_BYTE(SHR(udiSessionExpiry, 16) AND 16#FF);    uiPos := uiPos + 1;
        aTxBuf[uiPos] := UDINT_TO_BYTE(SHR(udiSessionExpiry, 8) AND 16#FF); uiPos := uiPos + 1;
        aTxBuf[uiPos] := UDINT_TO_BYTE(udiSessionExpiry AND 16#FF);         uiPos := uiPos + 1;
    END_IF

    // Receive Maximum
    aTxBuf[uiPos] := GVL_Mqtt.cnPropReceiveMaximum; uiPos := uiPos + 1;
    aTxBuf[uiPos] := UINT_TO_BYTE(SHR(uiReceiveMax, 8));    uiPos := uiPos + 1;
    aTxBuf[uiPos] := UINT_TO_BYTE(uiReceiveMax AND 16#FF);  uiPos := uiPos + 1;

    // Maximum Packet Size
    aTxBuf[uiPos] := GVL_Mqtt.cnPropMaxPacketSize;  uiPos := uiPos + 1;
    aTxBuf[uiPos] := UDINT_TO_BYTE(SHR(udMaxPacketSize, 24) AND 16#FF); uiPos := uiPos + 1;
    aTxBuf[uiPos] := UDINT_TO_BYTE(SHR(udMaxPacketSize, 16) AND 16#FF); uiPos := uiPos + 1;
    aTxBuf[uiPos] := UDINT_TO_BYTE(SHR(udMaxPacketSize, 8) AND 16#FF);  uiPos := uiPos + 1;
    aTxBuf[uiPos] := UDINT_TO_BYTE(udMaxPacketSize AND 16#FF);          uiPos := uiPos + 1;

    // Topic Alias Maximum
    uiTopicAliasMax := GVL_Mqtt.cnMaxTopicAlias;
    aTxBuf[uiPos] := GVL_Mqtt.cnPropTopicAliasMax;  uiPos := uiPos + 1;
    aTxBuf[uiPos] := UINT_TO_BYTE(SHR(uiTopicAliasMax, 8)); uiPos := uiPos + 1;
    aTxBuf[uiPos] := UINT_TO_BYTE(uiTopicAliasMax AND 16#FF); uiPos := uiPos + 1;

    // Request Problem Information
    aTxBuf[uiPos] := GVL_Mqtt.cnPropRequestProblemInfo; uiPos := uiPos + 1;
    aTxBuf[uiPos] := BOOL_TO_BYTE(bRequestProblemInfo);             uiPos := uiPos + 1;

    // Request Response Information
    IF bRequestResponseInfo THEN
        aTxBuf[uiPos] := GVL_Mqtt.cnPropRequestResponseInfo; uiPos := uiPos + 1;
        aTxBuf[uiPos] := 1; uiPos := uiPos + 1;
    END_IF
END_IF

// ****************** 载荷 = 客户端标识符 + 遗嘱主题 + 遗嘱消息 + 用户名 + 密码 ******************
// 客户端标识符
uiPos := uiPos + M_AppendString(sClientID, ADR(aTxBuf[uiPos]));

// 遗嘱主题 + 遗嘱消息
IF bWillFlag THEN
    IF eVersion = E_MqttVersion.byMqttVersion50 THEN
        aTxBuf[uiPos] := 0;
        uiPos := uiPos + 1;
    END_IF
    uiPos := uiPos + M_AppendString(sWillTopic, ADR(aTxBuf[uiPos]));
    uiPos := uiPos + M_AppendString(sWillMessage, ADR(aTxBuf[uiPos]));
END_IF

// 用户名
IF sUsername <> '' THEN
    uiPos := uiPos + M_AppendString(sUsername, ADR(aTxBuf[uiPos]));
END_IF

// 密码
IF sPassword <> '' THEN
    uiPos := uiPos + M_AppendString(sPassword, ADR(aTxBuf[uiPos]));
END_IF

uiTxLength := uiPos;                                                            //需要发送的报文总字节数
M_BuildConnectPacket := TRUE;

完整代码 2:M_HandleConnAck

  • 对应源码路径:10 MQTT/MqttClient_V1_0/Device/Application/MQTT/POUs/MqttClient NBS/FB_MqttClient/处理接收报文/M_HandleConnAck.st
  • 复制使用说明:这是 CONNECT 对应的收包处理代码,负责把 Broker 回来的 CONNACK 解析成状态机可用的数据。
  • 阅读重点:先看 Reason Code 判定,再看 MQTT 5.0 属性遍历,最后看 xWaitingForAck、M_InflightClear()、M_SubListClear() 这些连接收尾动作。
iecst
/// =======================================================================
/// 名称      : M_HandleConnAck
/// 功能      : 处理 CONNACK 接收报文
/// 说明      : 解析连接确认结果与服务器属性,成功后清除等待确认状态。
/// 编程人员  : ControlRookie
/// 时间      : 2026-05-05
/// 版本      : V1.1
/// =======================================================================
{attribute 'hide_all_locals'}
METHOD M_HandleConnAck : BOOL
VAR
    byReasonCode            : BYTE;  // CONNACK 原因码
    byConnAckFlags          : BYTE;  // CONNACK 标志位
    uiPropsLen              : UINT;  // 属性总长度
    uiPropertyHeaderBytes   : UINT;  // 属性长度字段字节数
    uiPropsStart            : UINT;  // 属性起始偏移
    uiPropsEnd              : UINT;  // 属性结束偏移
    uiPos                   : UINT;  // 当前解析偏移
    byPropId                : BYTE;  // 当前属性标识符
    uiLoopGuard             : UINT;  // 循环安全计数器
END_VAR

// === IMPLEMENTATION ===
IF uiRxLength < 4 THEN
    M_HandleConnAck := FALSE;
    RETURN;
END_IF

IF NOT xWaitingForAck OR aRxBuf[0] <> byExpectedMsgType THEN
    M_HandleConnAck := FALSE;
    RETURN;
END_IF

byConnAckFlags := aRxBuf[2];

IF (byConnAckFlags AND 16#FE) <> 0 THEN
    sDiagMsg := 'Invalid CONNACK flags';
    M_HandleConnAck := FALSE;
    RETURN;
END_IF

IF eVersion = E_MqttVersion.byMqttVersion50 THEN
    byReasonCode := aRxBuf[3];
ELSE
    byReasonCode := aRxBuf[3];
END_IF

IF byReasonCode <> 0 THEN
    IF eVersion = E_MqttVersion.byMqttVersion50 THEN
        CASE byReasonCode OF
            E_MqttReasonCode.byUnsupportedProtocolVersion: sDiagMsg := 'Protocol version not supported';
            E_MqttReasonCode.byClientIdentifierInvalid: sDiagMsg := 'Client identifier invalid';
            E_MqttReasonCode.byBadUserNamePassword: sDiagMsg := 'Bad username or password';
            E_MqttReasonCode.byNotAuthorized: sDiagMsg := 'Not authorized';
            E_MqttReasonCode.byServerUnavailable: sDiagMsg := 'Server unavailable';
            E_MqttReasonCode.byServerBusy: sDiagMsg := 'Server busy';
            E_MqttReasonCode.byBanned: sDiagMsg := 'Banned';
            E_MqttReasonCode.byServerShuttingDown: sDiagMsg := 'Server shutting down';
            E_MqttReasonCode.byBadAuthMethod: sDiagMsg := 'Bad authentication method';
            E_MqttReasonCode.byQuotaExceeded: sDiagMsg := 'Quota exceeded';
            E_MqttReasonCode.byRetainNotSupported: sDiagMsg := 'Retain not supported';
            E_MqttReasonCode.byQosNotSupported: sDiagMsg := 'QoS not supported';
            E_MqttReasonCode.byUseAnotherServer: sDiagMsg := 'Use another server';
            E_MqttReasonCode.byServerMoved: sDiagMsg := 'Server moved';
        ELSE
            sDiagMsg := CONCAT('CONNACK reason code: ', BYTE_TO_STRING(byReasonCode));
        END_CASE
    ELSE
        CASE byReasonCode OF
            1: sDiagMsg := 'Unacceptable protocol version';
            2: sDiagMsg := 'Identifier rejected';
            3: sDiagMsg := 'Server unavailable';
            4: sDiagMsg := 'Bad username or password';
            5: sDiagMsg := 'Not authorized';
        ELSE
            sDiagMsg := CONCAT('CONNACK return code: ', BYTE_TO_STRING(byReasonCode));
        END_CASE
    END_IF
    M_HandleConnAck := FALSE;
    RETURN;
END_IF

uiServerReceiveMax := GVL_Mqtt.cnDefaultReceiveMax;
byServerMaxQoS := TO_BYTE(E_MqttQoS.byQoS2);
bServerRetainAvailable := TRUE;
udServerMaxPacketSize := GVL_Mqtt.cnMaxPacketSize;
uiServerTopicAliasMax := 0;
uiSendQuota := GVL_Mqtt.cnDefaultReceiveMax;
bServerWildcardSubAvail := TRUE;
bServerSubIdAvail := TRUE;
bServerSharedSubAvail := TRUE;

IF eVersion = E_MqttVersion.byMqttVersion50 AND uiRxLength >= 5 THEN
    uiPos := 4;
    uiPropertyHeaderBytes := M_DecodeRemainingLength(pBuffer := ADR(aRxBuf[uiPos]), uiLength => uiPropsLen);
    IF uiPropertyHeaderBytes = 0 THEN
        M_HandleConnAck := FALSE;
        RETURN;
    END_IF
    uiPos := uiPos + uiPropertyHeaderBytes;
    uiPropsStart := uiPos;
    uiPropsEnd := uiPropsStart + uiPropsLen;
    IF uiPropsEnd > uiRxLength THEN
        M_HandleConnAck := FALSE;
        RETURN;
    END_IF
    uiLoopGuard := 0;

    WHILE uiPos < uiPropsEnd AND uiPos < uiRxLength AND uiLoopGuard < GVL_Mqtt.cnMaxPropertyLoop DO
        byPropId := aRxBuf[uiPos];
        uiPos := uiPos + 1;

        CASE byPropId OF
            GVL_Mqtt.cnPropReceiveMaximum:
                IF uiPos + 1 >= uiPropsEnd THEN
                    M_HandleConnAck := FALSE;
                    RETURN;
                END_IF
                uiServerReceiveMax := SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]);
                IF uiServerReceiveMax = 0 THEN
                    sDiagMsg := 'Receive Maximum must not be zero';
                    M_HandleConnAck := FALSE;
                    RETURN;
                END_IF
                uiPos := uiPos + 2;

            GVL_Mqtt.cnPropMaxQoS:
                IF uiPos >= uiPropsEnd THEN
                    M_HandleConnAck := FALSE;
                    RETURN;
                END_IF
                byServerMaxQoS := aRxBuf[uiPos];
                IF byServerMaxQoS > TO_BYTE(E_MqttQoS.byQoS2) THEN
                    sDiagMsg := 'Server Max QoS is invalid';
                    M_HandleConnAck := FALSE;
                    RETURN;
                END_IF
                uiPos := uiPos + 1;

            GVL_Mqtt.cnPropRetainAvailable:
                IF uiPos >= uiPropsEnd THEN
                    M_HandleConnAck := FALSE;
                    RETURN;
                END_IF
                bServerRetainAvailable := BYTE_TO_BOOL(aRxBuf[uiPos]);
                uiPos := uiPos + 1;

            GVL_Mqtt.cnPropMaxPacketSize:
                IF uiPos + 3 >= uiPropsEnd THEN
                    M_HandleConnAck := FALSE;
                    RETURN;
                END_IF
                udServerMaxPacketSize := SHL(BYTE_TO_UDINT(aRxBuf[uiPos]), 24) OR SHL(BYTE_TO_UDINT(aRxBuf[uiPos + 1]), 16) OR SHL(BYTE_TO_UDINT(aRxBuf[uiPos + 2]), 8) OR BYTE_TO_UDINT(aRxBuf[uiPos + 3]);
                IF udServerMaxPacketSize = 0 THEN
                    sDiagMsg := 'Server Maximum Packet Size must not be zero';
                    M_HandleConnAck := FALSE;
                    RETURN;
                END_IF
                uiPos := uiPos + 4;

            GVL_Mqtt.cnPropTopicAliasMax:
                IF uiPos + 1 >= uiPropsEnd THEN
                    M_HandleConnAck := FALSE;
                    RETURN;
                END_IF
                uiServerTopicAliasMax := SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]);
                uiPos := uiPos + 2;

            GVL_Mqtt.cnPropReasonString:
                IF uiPos + 1 >= uiPropsEnd THEN
                    M_HandleConnAck := FALSE;
                    RETURN;
                END_IF
                uiPos := uiPos + 2 + (SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]));

            GVL_Mqtt.cnPropServerKeepAlive:
                IF uiPos + 1 >= uiPropsEnd THEN
                    M_HandleConnAck := FALSE;
                    RETURN;
                END_IF
                uiKeepAlive := SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]);
                uiPos := uiPos + 2;

            GVL_Mqtt.cnPropAssignedClientId:
                IF uiPos + 1 >= uiPropsEnd THEN
                    M_HandleConnAck := FALSE;
                    RETURN;
                END_IF
                uiPos := uiPos + 2 + (SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]));

            GVL_Mqtt.cnPropAuthMethod,
            GVL_Mqtt.cnPropAuthData,
            GVL_Mqtt.cnPropServerReference:
                IF uiPos + 1 >= uiPropsEnd THEN
                    M_HandleConnAck := FALSE;
                    RETURN;
                END_IF
                uiPos := uiPos + 2 + (SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]));

            GVL_Mqtt.cnPropResponseInfo:
                IF uiPos + 1 >= uiPropsEnd THEN
                    M_HandleConnAck := FALSE;
                    RETURN;
                END_IF
                uiPos := uiPos + 2 + (SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]));

            GVL_Mqtt.cnPropUserProperty:
                IF uiPos + 1 >= uiPropsEnd THEN
                    M_HandleConnAck := FALSE;
                    RETURN;
                END_IF
                uiPos := uiPos + 2 + (SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]));
                IF uiPos + 1 >= uiPropsEnd THEN
                    M_HandleConnAck := FALSE;
                    RETURN;
                END_IF
                uiPos := uiPos + 2 + (SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]));

            GVL_Mqtt.cnPropWildcardSubAvail,
            GVL_Mqtt.cnPropSubIdAvail,
            GVL_Mqtt.cnPropSharedSubAvail:
                IF uiPos >= uiPropsEnd THEN
                    M_HandleConnAck := FALSE;
                    RETURN;
                END_IF
                CASE byPropId OF
                    GVL_Mqtt.cnPropWildcardSubAvail:
                        bServerWildcardSubAvail := BYTE_TO_BOOL(aRxBuf[uiPos]);
                    GVL_Mqtt.cnPropSubIdAvail:
                        bServerSubIdAvail := BYTE_TO_BOOL(aRxBuf[uiPos]);
                    GVL_Mqtt.cnPropSharedSubAvail:
                        bServerSharedSubAvail := BYTE_TO_BOOL(aRxBuf[uiPos]);
                END_CASE
                uiPos := uiPos + 1;

            GVL_Mqtt.cnPropRequestProblemInfo,
            GVL_Mqtt.cnPropRequestResponseInfo:
                IF uiPos >= uiPropsEnd THEN
                    M_HandleConnAck := FALSE;
                    RETURN;
                END_IF
                uiPos := uiPos + 1;

        ELSE
            M_HandleConnAck := FALSE;
            RETURN;
        END_CASE

        uiLoopGuard := uiLoopGuard + 1;
    END_WHILE
END_IF

uiRxLength := 0;
xWaitingForAck := FALSE;

IF eVersion = E_MqttVersion.byMqttVersion50 THEN
    uiSendQuota := uiServerReceiveMax;
END_IF

IF NOT bCleanSession THEN
    IF (byConnAckFlags AND GVL_Mqtt.cnConnAckSessionPresent) = 0 THEN
        M_InflightClear();
        M_SubListClear();
    END_IF
ELSE
    M_InflightClear();
    M_SubListClear();
END_IF

M_TopicAliasClear();

M_HandleConnAck := TRUE;

系列导航

  • 系列定位:第 2 篇
  • 上一篇:第1篇 MQTT 到底跑在哪一层
  • 下一篇:第3篇 PUBLISH 报文怎么写,QoS0 / QoS1 / QoS2 到底差在哪
评论和回复区

评论区预留

这里先保留评论和回复结构,不接入第三方服务。后续统一决定登录、匿名、审核、反垃圾和静态站兼容策略。

↑ ↓