ControlRookie
返回文章

第4篇_为什么 QoS1 / QoS2 最容易掉线超时?把 ACK、重发和状态机讲透

这一篇专门讲高可靠链路为什么最容易出问题,核心内容包括 ACK 链、inflight、超时重发、DUP 位以及高频场景下为什么状态机会最先暴露短板。

这一篇专门讲高可靠链路为什么最容易出问题,核心内容包括 ACK 链、inflight、超时重发、DUP 位以及高频场景下为什么状态机会最先暴露短板。

适合谁收藏

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

只要 MQTT 客户端一进现场,最常见的一类评价通常是这样的:

  • QoS0 正常
  • 一上 QoS1 就开始偶发 timeout
  • 一上 QoS2 更容易暴露问题
  • 高频时更容易掉进断线重连

这不是巧合。 也不是 MQTT 故意折磨人。

真正原因就一句话:

QoS1 / QoS2 不只是“消息等级更高”,而是要求客户端同时把 ACK、状态机、inflight、超时、重发 这几件事一起干对。

这一篇我们就把这条链讲透。


一、先给工程结论

为什么 QoS1 / QoS2 最容易炸?

因为它们不是“发一个包”的问题,而是“发完之后还要维持协议上下文”的问题。

你只要有任意一层没闭环:

  • ACK 匹配错
  • ACK 处理慢
  • 等待状态释放不及时
  • inflight 记录没清掉
  • 超时重发没接上

最后表现出来的,就很可能是:

  • Immediate response timeout
  • PubAck timeout
  • PubRec timeout
  • PubComp timeout
  • 状态机跳到 iTcpDisconnect

二、先把 QoS1 链路压成一张图

Mermaid
sequenceDiagram
    participant PLC as PLC Client
    participant Broker as Broker

    PLC->>Broker: PUBLISH(QoS1)
    Broker-->>PLC: PUBACK

表面看只有两步。 但客户端内部实际至少要做下面这些事:

  1. 构造 PUBLISH
  2. 分配 Packet ID
  3. 建立 inflight
  4. 发送报文
  5. 标记“我正在等 PUBACK”
  6. 收到 PUBACK 后匹配 Packet ID
  7. 删除 inflight
  8. 释放等待状态
  9. 回到 iConnected

所以别被“图很短”骗了。 短的是网络报文,不短的是客户端内部控制逻辑。


三、QoS2 链路为什么更考验实现水平

QoS2 完整链路是:

Mermaid
sequenceDiagram
    participant PLC as PLC Client
    participant Broker as Broker

    PLC->>Broker: PUBLISH(QoS2)
    Broker-->>PLC: PUBREC
    PLC->>Broker: PUBREL
    Broker-->>PLC: PUBCOMP

它的工程难点在于:

  • 第一段确认:消息到了
  • 第二段确认:这条消息可以安全完成,且只完成一次

所以 QoS2 本质上不是“QoS1 再加一个 ACK”。 它更像一条 双段状态链。


四、把 QoS1 / QoS2 放回状态机里看

Mermaid
flowchart TD
    A[iConnected] --> B[iPublish]
    B --> C[QoS0: iConnected]
    B --> D[QoS1: iPubAck]
    B --> E[QoS2: iPubRec]
    E --> F[iPubRel]
    F --> G[iPubComp]
    D --> A
    G --> A

这张图很重要。 因为它说明一个事实:

MQTT 高可靠发布,不是由某一个方法单独完成的,而是由整条状态链共同完成的。

所以你现场一旦看到:

  • 卡在 iPubAck
  • 卡在 iPubRec
  • 卡在 iPubComp

就不要再抽象地说“QoS 有问题”。 你应该立刻问:

  • 我现在在等哪一跳 ACK?
  • 这跳 ACK 有没有收到?
  • 收到了有没有被接收路径及时消化?
  • inflight 状态有没有推进?

五、真正支撑 QoS 的核心:inflight

这一层必须单独讲。

1. inflight 是什么

你可以先把 inflight 理解成:

客户端对“还没正式完成的 QoS 消息”的一本台账。

2. 为什么必须有它

因为 PLC 不是脚本环境,不会“一口气执行完所有网络逻辑”。 它是在扫描周期里一轮一轮推进状态的。

所以只要消息发出去之后还要继续等待:

  • PUBACK
  • PUBREC
  • PUBCOMP

你就必须把上下文保存起来。

3. 它至少保存什么

字段作用
Packet ID匹配 ACK
QoS决定走哪条确认链
Topic / Payload重发时复用
Retain重发时保持一致
tLastSend判断是否超时
uiRetryCount判断是否超过重试上限
bDup重发时置位
当前在途状态判断回到 iPublish 还是 iPubRel

六、超时重发到底怎么实现

这个库里,重发扫描是靠 M_InflightCheckTimeout 做的。

真实代码非常值得看:

iecst
FOR i := 1 TO GVL_Mqtt.cnMaxInflight DO
    IF aInflight[i].bUsed THEN
        IF (tCurrent - aInflight[i].tLastSend) >= GVL_Mqtt.cnInflightTimeout THEN
            IF aInflight[i].uiRetryCount < GVL_Mqtt.cnMaxRetries THEN
                aInflight[i].bDup := TRUE;
                aInflight[i].uiRetryCount := aInflight[i].uiRetryCount + 1;
                aInflight[i].tLastSend := tCurrent;
                M_InflightCheckTimeout := i;
                RETURN;
            ELSE
                M_SetError(
                    uiErrorCode := TO_UINT(E_ReasonCode.uiErrTimeout),
                    sMessage := 'Inflight publish retry exceeded');
                M_InflightRemove(uiPacketId := aInflight[i].uiPacketId);
            END_IF
        END_IF
    END_IF
END_FOR

这段代码说明了 4 件事:

  1. 它不是盲目全量重发,而是扫描首个超时项
  2. 超时后会把 bDup 置位
  3. 会累加重试次数
  4. 超过上限后,不再死循环,而是报错并清理

这就是为什么“能不能超时重发”不是一句空话,而是具体状态管理能力。


七、DUP 位到底什么时候才有意义

很多人知道 PUBLISH 首字节里有个 DUP 位。 但不知道它什么时候该置位。

结论很简单:

第一次发,不该乱置 DUP。

这个库里就是靠 inflight 超时后改 bDup 来驱动的。

也就是说,DUP 不是“高 QoS 就自动为 1”。 它表达的是:

这条消息不是第一次发,而是我重新发的。

八、为什么高频场景更容易暴露问题

高频时,客户端同时在干的事更多:

  • 处理新的发布请求
  • 等待上一条消息的 ACK
  • 处理入站 ACK 报文
  • 可能还要立即回对端 ACK
  • 还要扫 inflight 超时
  • KeepAlive 也可能同时到期

所以高频场景本质上是在压测两件事:

  1. 状态机是不是清晰
  2. 接收面和发送面有没有互相堵住

如果这两件事任何一个处理不好,就很容易出现一种现象:

明明网络没断、消息偶尔也收到了,但状态机却直接走到了 iTcpDisconnect。

这通常不是“网炸了”,而是协议时序没跟上。


九、为什么接收面即时消化这么关键

这个库有个很关键的方法:M_ProcessPendingFrames。

它的核心逻辑很短:

iecst
WHILE uiRxLength > 0 DO
    IF NOT M_ProcessReceive() THEN
        EXIT;
    END_IF
    bProcessed := TRUE;

    IF xPendingImmediateTx THEN
        EXIT;
    END_IF
END_WHILE

这段短代码,其实体现了一个很成熟的思路:

只要缓冲区里还有完整帧,就持续处理。

这件事特别重要。 因为 QoS1 / QoS2 场景下,很多掉线不是因为“没收到包”,而是因为:

  • 包收到了
  • 但是 ACK 没有及时发出去
  • 对端等不到,就按超时处理

十、为什么 Immediate response timeout 这么常见

这个报错名字本身已经把问题说了一半:

对某些入站报文,你不是“以后找机会再回”,而是必须尽快回。

典型场景包括:

  • 对端发来 QoS1 的 PUBLISH,你要回 PUBACK
  • 对端发来 QoS2 的 PUBLISH,你要回 PUBREC
  • 对端发来 PUBREL,你要回 PUBCOMP

如果这类即时响应被主状态机阻塞住,或者接收面没有及时把帧处理掉,就很容易出现:

  • 用户看到消息确实收到了
  • 但状态机还是报 Immediate response timeout

这类问题比“完全收不到消息”更隐蔽,也更像现场会遇到的真问题。


十一、怎么判断问题到底卡在哪一层

现场排这种问题时,建议直接按下面这个表走。

现象第一怀疑点
QoS0 正常,QoS1 不稳PUBACK 等待与 inflight 清理
QoS1 正常,QoS2 不稳PUBREC/PUBREL/PUBCOMP 状态推进
收到消息但仍报 Immediate timeout接收面即时 ACK 没及时发
高频时更容易炸iConnected 调度能力不足
一超时就掉进 iTcpDisconnect超时策略直接走统一异常收口

这比泛泛地说“看看网络吧”有效得多。


十二、把这条链完整串一遍

1. 标准层

  • QoS1:至少一次
  • QoS2:只一次

2. 报文层

  • QoS1:PUBLISH -> PUBACK
  • QoS2:PUBLISH -> PUBREC -> PUBREL -> PUBCOMP

3. 状态机层

  • iPublish
  • iPubAck
  • iPubRec
  • iPubRel
  • iPubComp

4. 实现层

能力方法 / 机制
构包M_BuildPublishPacket
超时扫描M_InflightCheckTimeout
接收帧循环消化M_ProcessPendingFrames
真正解析 ACK / PUBLISHM_ProcessReceive

十三、这一篇你最该记住的 6 句话

  1. QoS1 / QoS2 的难点从来不是“发包”,而是“发出去以后怎么把后面的链走完”。
  2. 高可靠 MQTT 客户端必须维护 inflight,而不是只靠当前输入参数。
  3. 超时重发不只是再发一次,而是要带着状态、次数、DUP 位一起推进。
  4. 高频场景测的不是网速,测的是状态机调度能力。
  5. 很多 Immediate response timeout,不是没收到消息,而是没及时把协议 ACK 发回去。
  6. QoS1 / QoS2 能长期跑稳,才说明客户端真正成熟。

十四、下篇预告

下一篇我们转到另一块非常容易写错的内容:

SUBSCRIBE / UNSUBSCRIBE 和主题过滤器校验

重点会讲:

  • #、+、$ 这些规则到底怎么落地
  • 共享订阅为什么不是随便写个 $share/... 就完了
  • 为什么取消订阅也不能只会“发个包”

从下一篇开始,工程味会继续上来。 因为订阅链一旦写错,现场通常不是“完全不工作”,而是“偶尔能用,但总有奇怪边界条件”。


完整 ST 代码

复制使用说明

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

代码阅读重点

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

完整代码 1:M_ProcessReceive

  • 对应源码路径:10 MQTT/MqttClient_V1_0/Device/Application/MQTT/POUs/MqttClient NBS/FB_MqttClient/处理接收报文/M_ProcessReceive.st
  • 复制使用说明:这是 QoS1/QoS2 最关键的收包分发入口,所有 ACK 类报文几乎都会经过这里。
  • 阅读重点:先看如何识别报文类型,再看不同 ACK 被分派给哪个 M_Handle... 方法,最后看它对 xWaitingForAck、xPendingImmediateTx 的影响。
iecst
/// =======================================================================
/// 名称      : M_ProcessReceive
/// 功能      : 统一处理接收到的 MQTT 报文
/// 说明      : 校验完整帧后,按报文类型执行解析、确认和状态更新。
/// =======================================================================
{attribute 'hide_all_locals'}
METHOD M_ProcessReceive : BOOL
VAR
    byHeader             : BYTE;          // 固定报头首字节
    byMsgType            : BYTE;          // 报文类型
    byFlags              : BYTE;          // 固定报头标志
    byReasonCode         : BYTE;          // 原因码
    bySubAckCode         : BYTE;          // SUBACK 原因码
    uiFixedHeaderLen     : UINT;          // 固定报头总长度
    uiRemainingLen       : UINT;          // 剩余长度
    uiFrameEndPos        : UINT;          // 当前帧结束位置
    uiTopicPos           : UINT;          // 主题起始偏移
    uiTopicLen           : UINT;          // 主题长度
    uiPayloadPos         : UINT;          // 载荷起始偏移
    uiPayloadLen         : UINT;          // 载荷长度
    uiPacketId           : UINT;          // 报文标识符
    uiAliasId            : UINT;          // 主题别名
    udiSubIdentifier     : UDINT;         // 订阅标识符
    uiPropsLen           : UINT;          // 属性长度
    uiPropsHeaderLen     : UINT;          // 属性长度字段字节数
    uiPropsEnd           : UINT;          // 属性区结束偏移
    uiPropsPos           : UINT;          // 属性读取偏移
    uiIndex              : UINT;          // 通用索引
    uiInflightIndex      : UINT;          // 在途索引
    uiRxQoS2Index        : UINT;          // QoS2 去重索引
    i                    : UINT;          // 循环索引
    j                    : DINT;          // 历史数组循环索引
    sAliasTopic          : STRING(GVL_Mqtt.cnMaxTopicLen); // 别名映射主题
END_VAR

// === IMPLEMENTATION ===
IF uiRxLength < 2 THEN
    M_ProcessReceive := FALSE;
    RETURN;
END_IF

byHeader := aRxBuf[0];
byMsgType := byHeader AND GVL_Mqtt.cnHdrTypeMask;
byFlags := byHeader AND GVL_Mqtt.cnHdrFlagsMask;

uiFixedHeaderLen := 1 + M_DecodeRemainingLength(
    pBuffer := ADR(aRxBuf[1]),
    uiLength => uiRemainingLen);

IF uiFixedHeaderLen = 0 THEN
    M_ProcessReceive := FALSE;
    RETURN;
END_IF

uiFrameEndPos := uiFixedHeaderLen + uiRemainingLen;
IF uiRxLength < uiFrameEndPos THEN
    M_ProcessReceive := FALSE;
    RETURN;
END_IF

CASE byMsgType OF
    E_MqttPacketType.byConnAck:
        M_SetError(
            uiErrorCode := TO_UINT(E_ReasonCode.uiErrProtocolError),
            sMessage := 'Unexpected CONNACK packet');
        eState := E_MqttState.iTcpDisconnect;
        M_ProcessReceive := FALSE;
        RETURN;

    E_MqttPacketType.byPublish:
        byReceivedQoS := SHR(byFlags AND GVL_Mqtt.cnHdrQoSMask, 1);
        uiAliasId := 0;
        udiSubIdentifier := 0;
        bReceivedRetain := (byFlags AND GVL_Mqtt.cnHdrRetainFlag) <> 0;

        IF (byReceivedQoS > TO_BYTE(E_MqttQoS.byQoS2)) OR
           ((byReceivedQoS = TO_BYTE(E_MqttQoS.byQoS0)) AND ((byFlags AND GVL_Mqtt.cnHdrDupFlag) <> 0)) THEN
            M_SetError(
                uiErrorCode := TO_UINT(E_ReasonCode.uiErrProtocolError),
                sMessage := 'Invalid PUBLISH header flags');
            eState := E_MqttState.iTcpDisconnect;
            M_ProcessReceive := FALSE;
            RETURN;
        END_IF

        uiTopicPos := uiFixedHeaderLen;
        IF uiTopicPos + 1 >= uiFrameEndPos THEN
            M_ProcessReceive := FALSE;
            RETURN;
        END_IF

        uiTopicLen := SHL(BYTE_TO_UINT(aRxBuf[uiTopicPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiTopicPos + 1]);
        uiTopicPos := uiTopicPos + 2;
        uiPayloadPos := uiTopicPos + uiTopicLen;

        sRecTopic := '';
        IF uiTopicLen > 0 THEN
            IF uiPayloadPos > uiFrameEndPos THEN
                M_ProcessReceive := FALSE;
                RETURN;
            END_IF
            FOR i := 1 TO uiTopicLen DO
                sRecTopic := CONCAT(sRecTopic, Util.WORD_AS_STRING(aRxBuf[uiTopicPos + i - 1], FALSE));
            END_FOR
        END_IF

        uiPacketId := 0;
        IF byReceivedQoS > TO_BYTE(E_MqttQoS.byQoS0) THEN
            IF uiPayloadPos + 1 >= uiFrameEndPos THEN
                M_ProcessReceive := FALSE;
                RETURN;
            END_IF
            uiPacketId := SHL(BYTE_TO_UINT(aRxBuf[uiPayloadPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPayloadPos + 1]);
            IF uiPacketId = 0 THEN
                M_SetError(
                    uiErrorCode := TO_UINT(E_ReasonCode.uiErrProtocolError),
                    sMessage := 'Packet ID must not be zero');
                eState := E_MqttState.iTcpDisconnect;
                M_ProcessReceive := FALSE;
                RETURN;
            END_IF
            uiPayloadPos := uiPayloadPos + 2;
        END_IF

        IF eVersion = E_MqttVersion.byMqttVersion50 THEN
            IF uiPayloadPos >= uiFrameEndPos THEN
                M_ProcessReceive := FALSE;
                RETURN;
            END_IF
            uiPropsHeaderLen := M_DecodeRemainingLength(
                pBuffer := ADR(aRxBuf[uiPayloadPos]),
                uiLength => uiPropsLen);
            IF uiPropsHeaderLen = 0 THEN
                M_ProcessReceive := FALSE;
                RETURN;
            END_IF
            uiPropsPos := uiPayloadPos + uiPropsHeaderLen;
            uiPropsEnd := uiPropsPos + uiPropsLen;
            IF uiPropsEnd > uiFrameEndPos THEN
                M_ProcessReceive := FALSE;
                RETURN;
            END_IF

            WHILE uiPropsPos < uiPropsEnd DO
                CASE aRxBuf[uiPropsPos] OF
                    GVL_Mqtt.cnPropTopicAlias:
                        uiPropsPos := uiPropsPos + 1;
                        IF uiPropsPos + 1 >= uiPropsEnd THEN
                            M_ProcessReceive := FALSE;
                            RETURN;
                        END_IF
                        uiAliasId := SHL(BYTE_TO_UINT(aRxBuf[uiPropsPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPropsPos + 1]);
                        uiPropsPos := uiPropsPos + 2;
                    GVL_Mqtt.cnPropSubscriptionId:
                        uiPropsPos := uiPropsPos + 1;
                        uiIndex := 0;
                        uiPropsHeaderLen := M_DecodeRemainingLength(
                            pBuffer := ADR(aRxBuf[uiPropsPos]),
                            uiLength => uiIndex);
                        IF uiPropsHeaderLen = 0 THEN
                            M_ProcessReceive := FALSE;
                            RETURN;
                        END_IF
                        udiSubIdentifier := TO_UDINT(uiIndex);
                        uiPropsPos := uiPropsPos + uiPropsHeaderLen;
                    GVL_Mqtt.cnPropPayloadFormat,
                    GVL_Mqtt.cnPropRetainAvailable,
                    GVL_Mqtt.cnPropRequestProblemInfo,
                    GVL_Mqtt.cnPropRequestResponseInfo,
                    GVL_Mqtt.cnPropWildcardSubAvail,
                    GVL_Mqtt.cnPropSubIdAvail,
                    GVL_Mqtt.cnPropSharedSubAvail:
                        uiPropsPos := uiPropsPos + 2;
                    GVL_Mqtt.cnPropReceiveMaximum,
                    GVL_Mqtt.cnPropTopicAliasMax,
                    GVL_Mqtt.cnPropServerKeepAlive:
                        uiPropsPos := uiPropsPos + 3;
                    GVL_Mqtt.cnPropMessageExpiry,
                    GVL_Mqtt.cnPropSessionExpiry,
                    GVL_Mqtt.cnPropMaxPacketSize:
                        uiPropsPos := uiPropsPos + 5;
                    GVL_Mqtt.cnPropContentType,
                    GVL_Mqtt.cnPropResponseTopic,
                    GVL_Mqtt.cnPropCorrelationData,
                    GVL_Mqtt.cnPropAssignedClientId,
                    GVL_Mqtt.cnPropAuthMethod,
                    GVL_Mqtt.cnPropAuthData,
                    GVL_Mqtt.cnPropResponseInfo,
                    GVL_Mqtt.cnPropServerReference,
                    GVL_Mqtt.cnPropReasonString:
                        uiPropsPos := uiPropsPos + 1;
                        IF uiPropsPos + 1 >= uiPropsEnd THEN
                            M_ProcessReceive := FALSE;
                            RETURN;
                        END_IF
                        uiPropsPos := uiPropsPos + 2 + (
                            SHL(BYTE_TO_UINT(aRxBuf[uiPropsPos]), 8) OR
                            BYTE_TO_UINT(aRxBuf[uiPropsPos + 1]));
                    GVL_Mqtt.cnPropUserProperty:
                        uiPropsPos := uiPropsPos + 1;
                        FOR i := 1 TO 2 DO
                            IF uiPropsPos + 1 >= uiPropsEnd THEN
                                M_ProcessReceive := FALSE;
                                RETURN;
                            END_IF
                            uiPropsPos := uiPropsPos + 2 + (
                                SHL(BYTE_TO_UINT(aRxBuf[uiPropsPos]), 8) OR
                                BYTE_TO_UINT(aRxBuf[uiPropsPos + 1]));
                        END_FOR
                ELSE
                    M_SetError(
                        uiErrorCode := TO_UINT(E_ReasonCode.uiErrProtocolError),
                        sMessage := 'Unsupported PUBLISH property');
                    eState := E_MqttState.iTcpDisconnect;
                    M_ProcessReceive := FALSE;
                    RETURN;
                END_CASE

                IF uiPropsPos > uiPropsEnd THEN
                    M_ProcessReceive := FALSE;
                    RETURN;
                END_IF
            END_WHILE

            uiPayloadPos := uiPropsEnd;

            IF uiAliasId > 0 THEN
                IF uiAliasId > GVL_Mqtt.cnMaxTopicAlias THEN
                    M_SetError(
                        uiErrorCode := TO_UINT(E_ReasonCode.uiErrTopicAliasInvalid),
                        sMessage := 'Topic Alias exceeds client maximum');
                    eState := E_MqttState.iTcpDisconnect;
                    M_ProcessReceive := FALSE;
                    RETURN;
                END_IF

                IF sRecTopic = '' THEN
                    sAliasTopic := M_TopicAliasLookup(uiAliasId := uiAliasId);
                    IF sAliasTopic = '' THEN
                        M_SetError(
                            uiErrorCode := TO_UINT(E_ReasonCode.uiErrTopicAliasInvalid),
                            sMessage := 'Unknown Topic Alias');
                        eState := E_MqttState.iTcpDisconnect;
                        M_ProcessReceive := FALSE;
                        RETURN;
                    END_IF
                    sRecTopic := sAliasTopic;
                ELSE
                    IF NOT M_TopicAliasRegister(uiAliasId := uiAliasId, sTopic := sRecTopic) THEN
                        M_SetError(
                            uiErrorCode := TO_UINT(E_ReasonCode.uiErrTopicAliasInvalid),
                            sMessage := 'Topic Alias registration failed');
                        eState := E_MqttState.iTcpDisconnect;
                        M_ProcessReceive := FALSE;
                        RETURN;
                    END_IF
                END_IF
            END_IF
        END_IF

        uiPayloadLen := uiFrameEndPos - uiPayloadPos;
        sRecPayload := '';
        IF uiPayloadLen > 0 THEN
            FOR i := 1 TO uiPayloadLen DO
                sRecPayload := CONCAT(sRecPayload, Util.WORD_AS_STRING(aRxBuf[uiPayloadPos + i - 1], FALSE));
            END_FOR
        END_IF

        IF byReceivedQoS = TO_BYTE(E_MqttQoS.byQoS1) THEN
            IF uiRxInFlightQosCount >= uiReceiveMax THEN
                M_SetError(
                    uiErrorCode := TO_UINT(E_ReasonCode.uiErrReceiveMaxExceeded),
                    sMessage := 'Receive Maximum exceeded');
                eState := E_MqttState.iTcpDisconnect;
                M_ProcessReceive := FALSE;
                RETURN;
            END_IF
            uiRxInFlightQosCount := uiRxInFlightQosCount + 1;
            IF NOT M_BuildPubAckPacket(uiPacketId := uiPacketId) THEN
                IF uiRxInFlightQosCount > 0 THEN
                    uiRxInFlightQosCount := uiRxInFlightQosCount - 1;
                END_IF
                M_SetError(
                    uiErrorCode := TO_UINT(E_ReasonCode.uiErrProtocolError),
                    sMessage := 'Build PUBACK failed');
                eState := E_MqttState.iTcpDisconnect;
                M_ProcessReceive := FALSE;
                RETURN;
            END_IF
            xPendingImmediateTx := TRUE;

            FOR j := UPPER_BOUND(aRecTopicList, 1) TO LOWER_BOUND(aRecTopicList, 1) + 1 BY -1 DO
                aRecTopicList[j] := aRecTopicList[j - 1];
                aRecPayloadList[j] := aRecPayloadList[j - 1];
            END_FOR
            aRecTopicList[0] := sRecTopic;
            aRecPayloadList[0] := sRecPayload;

            bMessageReceived := TRUE;
            uiMessagesReceived := uiMessagesReceived + 1;
            dtLastMessageTime := ULINT_TO_DT(uliSysTime / 1000);
        ELSIF byReceivedQoS = TO_BYTE(E_MqttQoS.byQoS2) THEN
            uiRxQoS2Index := M_RxQoS2Find(uiPacketId := uiPacketId);
            IF uiRxQoS2Index = 0 THEN
                IF uiRxInFlightQosCount >= uiReceiveMax THEN
                    M_SetError(
                        uiErrorCode := TO_UINT(E_ReasonCode.uiErrReceiveMaxExceeded),
                        sMessage := 'Receive Maximum exceeded');
                    eState := E_MqttState.iTcpDisconnect;
                    M_ProcessReceive := FALSE;
                    RETURN;
                END_IF
                IF NOT M_RxQoS2Add(uiPacketId := uiPacketId) THEN
                    M_SetError(
                        uiErrorCode := TO_UINT(E_ReasonCode.uiErrReceiveMaxExceeded),
                        sMessage := 'QoS2 dedup queue is full');
                    eState := E_MqttState.iTcpDisconnect;
                    M_ProcessReceive := FALSE;
                    RETURN;
                END_IF

                FOR j := UPPER_BOUND(aRecTopicList, 1) TO LOWER_BOUND(aRecTopicList, 1) + 1 BY -1 DO
                    aRecTopicList[j] := aRecTopicList[j - 1];
                    aRecPayloadList[j] := aRecPayloadList[j - 1];
                END_FOR
                aRecTopicList[0] := sRecTopic;
                aRecPayloadList[0] := sRecPayload;

                bMessageReceived := TRUE;
                uiMessagesReceived := uiMessagesReceived + 1;
                dtLastMessageTime := ULINT_TO_DT(uliSysTime / 1000);
            END_IF
            IF NOT M_BuildPubRecPacket(uiPacketId := uiPacketId) THEN
                M_SetError(
                    uiErrorCode := TO_UINT(E_ReasonCode.uiErrProtocolError),
                    sMessage := 'Build PUBREC failed');
                eState := E_MqttState.iTcpDisconnect;
                M_ProcessReceive := FALSE;
                RETURN;
            END_IF
            xPendingImmediateTx := TRUE;
        ELSE
            FOR j := UPPER_BOUND(aRecTopicList, 1) TO LOWER_BOUND(aRecTopicList, 1) + 1 BY -1 DO
                aRecTopicList[j] := aRecTopicList[j - 1];
                aRecPayloadList[j] := aRecPayloadList[j - 1];
            END_FOR
            aRecTopicList[0] := sRecTopic;
            aRecPayloadList[0] := sRecPayload;

            bMessageReceived := TRUE;
            uiMessagesReceived := uiMessagesReceived + 1;
            dtLastMessageTime := ULINT_TO_DT(uliSysTime / 1000);
        END_IF

    E_MqttPacketType.byPubAck:
        IF (byFlags <> 0) OR (uiFrameEndPos < 4) THEN
            M_ProcessReceive := FALSE;
            RETURN;
        END_IF
        uiPacketId := SHL(BYTE_TO_UINT(aRxBuf[uiFixedHeaderLen]), 8) OR BYTE_TO_UINT(aRxBuf[uiFixedHeaderLen + 1]);
        IF uiPacketId = uiExpectedPacketId THEN
            byReasonCode := 0;
            IF (eVersion = E_MqttVersion.byMqttVersion50) AND (uiFrameEndPos > uiFixedHeaderLen + 2) THEN
                byReasonCode := aRxBuf[uiFixedHeaderLen + 2];
            END_IF
            xWaitingForAck := FALSE;
            bDup := FALSE;
            IF byReasonCode >= 16#80 THEN
                M_SetError(
                    uiErrorCode := TO_UINT(E_ReasonCode.uiErrPubAckRefused),
                    sMessage := 'PUBACK rejected');
            ELSE
                xPublishedEvent := TRUE;
            END_IF
            M_InflightRemove(uiPacketId := uiPacketId);
        ELSE
            M_ProcessReceive := FALSE;
            RETURN;
        END_IF

    E_MqttPacketType.byPubRec:
        IF (byFlags <> 0) OR (uiFrameEndPos < 4) THEN
            M_ProcessReceive := FALSE;
            RETURN;
        END_IF
        uiPacketId := SHL(BYTE_TO_UINT(aRxBuf[uiFixedHeaderLen]), 8) OR BYTE_TO_UINT(aRxBuf[uiFixedHeaderLen + 1]);
        IF uiPacketId = uiExpectedPacketId THEN
            byReasonCode := 0;
            IF (eVersion = E_MqttVersion.byMqttVersion50) AND (uiFrameEndPos >= 5) THEN
                byReasonCode := aRxBuf[uiFixedHeaderLen + 2];
                IF byReasonCode >= 16#80 THEN
                    M_SetError(
                        uiErrorCode := TO_UINT(E_ReasonCode.uiErrPubRecRefused),
                        sMessage := 'PUBREC rejected');
                    M_InflightRemove(uiPacketId := uiPacketId);
                    xWaitingForAck := FALSE;
                END_IF
            END_IF
            IF NOT bError THEN
                M_InflightUpdateState(
                    uiPacketId := uiPacketId,
                    eNewState := E_MqttInflightState.iPubRecReceived);
                xWaitingForAck := FALSE;
                uiQoS2PacketId := uiPacketId;
            END_IF
        ELSE
            M_ProcessReceive := FALSE;
            RETURN;
        END_IF

    16#60:
        IF (byFlags <> 2) OR (uiFrameEndPos < 4) THEN
            M_SetError(
                uiErrorCode := TO_UINT(E_ReasonCode.uiErrProtocolError),
                sMessage := 'Invalid PUBREL flags');
            eState := E_MqttState.iTcpDisconnect;
            M_ProcessReceive := FALSE;
            RETURN;
        END_IF
        uiPacketId := SHL(BYTE_TO_UINT(aRxBuf[uiFixedHeaderLen]), 8) OR BYTE_TO_UINT(aRxBuf[uiFixedHeaderLen + 1]);
        M_RxQoS2Remove(uiPacketId := uiPacketId);
        IF NOT M_BuildPubCompPacket(uiPacketId := uiPacketId) THEN
            M_SetError(
                uiErrorCode := TO_UINT(E_ReasonCode.uiErrProtocolError),
                sMessage := 'Build PUBCOMP failed');
            eState := E_MqttState.iTcpDisconnect;
            M_ProcessReceive := FALSE;
            RETURN;
        END_IF
        xPendingImmediateTx := TRUE;

    E_MqttPacketType.byPubComp:
        IF (byFlags <> 0) OR (uiFrameEndPos < 4) THEN
            M_ProcessReceive := FALSE;
            RETURN;
        END_IF
        uiPacketId := SHL(BYTE_TO_UINT(aRxBuf[uiFixedHeaderLen]), 8) OR BYTE_TO_UINT(aRxBuf[uiFixedHeaderLen + 1]);
        IF uiPacketId = uiQoS2PacketId THEN
            byReasonCode := 0;
            IF (eVersion = E_MqttVersion.byMqttVersion50) AND (uiFrameEndPos > uiFixedHeaderLen + 2) THEN
                byReasonCode := aRxBuf[uiFixedHeaderLen + 2];
            END_IF
            xWaitingForAck := FALSE;
            bDup := FALSE;
            IF byReasonCode >= 16#80 THEN
                M_SetError(
                    uiErrorCode := TO_UINT(E_ReasonCode.uiErrPubCompRefused),
                    sMessage := 'PUBCOMP rejected');
            ELSE
                xPublishedEvent := TRUE;
            END_IF
            M_InflightRemove(uiPacketId := uiPacketId);
        ELSE
            M_ProcessReceive := FALSE;
            RETURN;
        END_IF

    E_MqttPacketType.bySubAck:
        IF (byFlags <> 0) OR (uiFrameEndPos < 5) THEN
            M_ProcessReceive := FALSE;
            RETURN;
        END_IF
        uiPacketId := SHL(BYTE_TO_UINT(aRxBuf[uiFixedHeaderLen]), 8) OR BYTE_TO_UINT(aRxBuf[uiFixedHeaderLen + 1]);
        IF uiPacketId = uiPendingSubPacketId THEN
            uiPayloadPos := uiFixedHeaderLen + 2;
            IF eVersion = E_MqttVersion.byMqttVersion50 THEN
                uiPropsHeaderLen := M_DecodeRemainingLength(
                    pBuffer := ADR(aRxBuf[uiPayloadPos]),
                    uiLength => uiPropsLen);
                IF uiPropsHeaderLen = 0 THEN
                    M_ProcessReceive := FALSE;
                    RETURN;
                END_IF
                IF uiPayloadPos + uiPropsHeaderLen + uiPropsLen > uiFrameEndPos THEN
                    M_ProcessReceive := FALSE;
                    RETURN;
                END_IF
                uiPayloadPos := uiPayloadPos + uiPropsHeaderLen + uiPropsLen;
            END_IF
            IF uiPayloadPos >= uiFrameEndPos THEN
                M_ProcessReceive := FALSE;
                RETURN;
            END_IF
            bySubAckCode := aRxBuf[uiPayloadPos];
            xWaitingForSubAck := FALSE;
            IF bySubAckCode < 16#80 THEN
                IF bySubAckCode > TO_BYTE(E_MqttQoS.byQoS2) THEN
                    M_SetError(
                        uiErrorCode := TO_UINT(E_ReasonCode.uiErrProtocolError),
                        sMessage := 'SUBACK reason code is invalid');
                    eState := E_MqttState.iTcpDisconnect;
                    M_ProcessReceive := FALSE;
                    RETURN;
                END_IF
                xSubscribedEvent := TRUE;
                M_SubListAdd(
                    sTopic := sSubTopic,
                    eQos := eSubQoS,
                    udiSubscriptionId := udiSubscriptionId);
            ELSE
                M_SetError(
                    uiErrorCode := TO_UINT(E_ReasonCode.uiErrSubscriptionFailed),
                    sMessage := 'Subscription rejected');
            END_IF
        ELSE
            M_ProcessReceive := FALSE;
            RETURN;
        END_IF

    E_MqttPacketType.byUnsubAck:
        IF byFlags <> 0 THEN
            M_ProcessReceive := FALSE;
            RETURN;
        END_IF
        IF ((eVersion = E_MqttVersion.byMqttVersion50) AND (uiFrameEndPos < 5)) OR
           ((eVersion <> E_MqttVersion.byMqttVersion50) AND (uiFrameEndPos < 4)) THEN
            M_ProcessReceive := FALSE;
            RETURN;
        END_IF
        uiPacketId := SHL(BYTE_TO_UINT(aRxBuf[uiFixedHeaderLen]), 8) OR BYTE_TO_UINT(aRxBuf[uiFixedHeaderLen + 1]);
        IF uiPacketId = uiPendingUnsubPacketId THEN
            xWaitingForUnsubAck := FALSE;
            byReasonCode := 0;
            IF eVersion = E_MqttVersion.byMqttVersion50 THEN
                uiPayloadPos := uiFixedHeaderLen + 2;
                uiPropsHeaderLen := M_DecodeRemainingLength(
                    pBuffer := ADR(aRxBuf[uiPayloadPos]),
                    uiLength => uiPropsLen);
                IF uiPropsHeaderLen = 0 THEN
                    M_ProcessReceive := FALSE;
                    RETURN;
                END_IF
                IF uiPayloadPos + uiPropsHeaderLen + uiPropsLen > uiFrameEndPos THEN
                    M_ProcessReceive := FALSE;
                    RETURN;
                END_IF
                uiPayloadPos := uiPayloadPos + uiPropsHeaderLen + uiPropsLen;
                IF uiPayloadPos < uiFrameEndPos THEN
                    byReasonCode := aRxBuf[uiPayloadPos];
                END_IF
            END_IF
            IF (byReasonCode <> 0) AND (byReasonCode <> GVL_Mqtt.cnRcNoSubscriptionExisted) AND (byReasonCode < 16#80) THEN
                M_SetError(
                    uiErrorCode := TO_UINT(E_ReasonCode.uiErrProtocolError),
                    sMessage := 'UNSUBACK reason code is invalid');
                eState := E_MqttState.iTcpDisconnect;
                M_ProcessReceive := FALSE;
                RETURN;
            END_IF
            IF byReasonCode >= 16#80 THEN
                M_SetError(
                    uiErrorCode := TO_UINT(E_ReasonCode.uiErrSubscriptionFailed),
                    sMessage := 'Unsubscribe rejected');
            ELSE
                xUnsubscribedEvent := TRUE;
                M_SubListRemove(sTopic := sUnsubTopic);
            END_IF
        ELSE
            M_ProcessReceive := FALSE;
            RETURN;
        END_IF

    E_MqttPacketType.byPingResp:
        IF (byFlags <> 0) OR (uiRemainingLen <> 0) THEN
            M_ProcessReceive := FALSE;
            RETURN;
        END_IF
        bPingPending := FALSE;

    E_MqttPacketType.byDisconnect:
        byReasonCode := 0;
        IF eVersion = E_MqttVersion.byMqttVersion50 THEN
            IF byFlags <> 0 THEN
                M_SetError(
                    uiErrorCode := TO_UINT(E_ReasonCode.uiErrProtocolError),
                    sMessage := 'Invalid DISCONNECT flags');
                eState := E_MqttState.iTcpDisconnect;
                M_ProcessReceive := FALSE;
                RETURN;
            END_IF
            IF uiRemainingLen >= 1 THEN
                byReasonCode := aRxBuf[uiFixedHeaderLen];
            END_IF
        ELSIF (byFlags <> 0) OR (uiRemainingLen <> 0) THEN
            M_SetError(
                uiErrorCode := TO_UINT(E_ReasonCode.uiErrProtocolError),
                sMessage := 'Invalid MQTT 3.1.1 DISCONNECT');
            eState := E_MqttState.iTcpDisconnect;
            M_ProcessReceive := FALSE;
            RETURN;
        END_IF
        M_SetError(
            uiErrorCode := TO_UINT(E_ReasonCode.uiErrNotConnected),
            sMessage := CONCAT('Server sent DISCONNECT rc=', BYTE_TO_STRING(byReasonCode)));
        eState := E_MqttState.iTcpDisconnect;

    E_MqttPacketType.byAuth:
        IF eVersion <> E_MqttVersion.byMqttVersion50 THEN
            M_SetError(
                uiErrorCode := TO_UINT(E_ReasonCode.uiErrProtocolError),
                sMessage := 'Unexpected AUTH packet');
            eState := E_MqttState.iTcpDisconnect;
        ELSE
            M_SetError(
                uiErrorCode := TO_UINT(E_ReasonCode.uiErrAuthFailed),
                sMessage := 'AUTH exchange is not supported in MqttClient_V1_0');
            eState := E_MqttState.iTcpDisconnect;
        END_IF

ELSE
        M_SetError(
            uiErrorCode := TO_UINT(E_ReasonCode.uiErrProtocolError),
            sMessage := 'Unsupported packet type');
        eState := E_MqttState.iTcpDisconnect;
        M_ProcessReceive := FALSE;
        RETURN;
END_CASE

IF uiRxLength > uiFrameEndPos THEN
    uiPayloadLen := uiRxLength - uiFrameEndPos;
    IF uiPayloadLen > 0 THEN
        FOR i := 1 TO uiPayloadLen DO
            aRxBuf[i - 1] := aRxBuf[uiFrameEndPos + i - 1];
        END_FOR
    END_IF
    uiRxLength := uiPayloadLen;
ELSE
    uiRxLength := 0;
END_IF

M_ProcessReceive := TRUE;

完整代码 2:M_InflightCheckTimeout

  • 对应源码路径:10 MQTT/MqttClient_V1_0/Device/Application/MQTT/POUs/MqttClient NBS/FB_MqttClient/辅助功能/M_InflightCheckTimeout.st
  • 复制使用说明:这个方法决定哪一条在途消息超时、是否该触发重发,是高频场景稳定性的关键检查点。
  • 阅读重点:重点看在途队列扫描、超时判定条件和返回索引,这一层直接决定状态机会不会进入重发分支。
iecst
/// =======================================================================
/// 名称      : M_InflightCheckTimeout
/// 功能      : 扫描在途消息超时
/// 说明      : 找到需要重发的首个在途消息并返回其索引
/// =======================================================================
{attribute 'hide_all_locals'}
METHOD M_InflightCheckTimeout : UINT
VAR
    i               : UINT;           // 循环索引
    tCurrent        : TIME;           // 当前时间
END_VAR

// === IMPLEMENTATION ===
tCurrent := TIME();

FOR i := 1 TO GVL_Mqtt.cnMaxInflight DO
    IF aInflight[i].bUsed THEN
        IF (tCurrent - aInflight[i].tLastSend) >= GVL_Mqtt.cnInflightTimeout THEN
            IF aInflight[i].uiRetryCount < GVL_Mqtt.cnMaxRetries THEN
                aInflight[i].bDup := TRUE;
                aInflight[i].uiRetryCount := aInflight[i].uiRetryCount + 1;
                aInflight[i].tLastSend := tCurrent;
                M_InflightCheckTimeout := i;
                RETURN;
            ELSE
                M_SetError(
                    uiErrorCode := TO_UINT(E_ReasonCode.uiErrTimeout),
                    sMessage := 'Inflight publish retry exceeded');
                M_InflightRemove(uiPacketId := aInflight[i].uiPacketId);
            END_IF
        END_IF
    END_IF
END_FOR

M_InflightCheckTimeout := 0;

完整代码 3:M_ProcessPendingFrames

  • 对应源码路径:10 MQTT/MqttClient_V1_0/Device/Application/MQTT/POUs/MqttClient NBS/FB_MqttClient/辅助功能/M_ProcessPendingFrames.st
  • 复制使用说明:这个方法负责把接收缓冲区里已经到达的完整帧持续吃掉,避免 QoS 高一点就积帧、超时、掉线。
  • 阅读重点:先看完整帧边界怎么判,再看循环消费逻辑,最后看为什么它能把“偶发来不及响应”变成“连续稳定处理”。
iecst
/// =======================================================================
/// 名称      : M_ProcessPendingFrames
/// 功能      : 连续处理缓冲区中的完整报文
/// 说明      : 只要存在完整报文就循环调用 M_ProcessReceive
/// =======================================================================
{attribute 'hide_all_locals'}
METHOD M_ProcessPendingFrames : BOOL
VAR
    bProcessed      : BOOL;           // 是否处理过至少一帧
END_VAR

// === IMPLEMENTATION ===
bProcessed := FALSE;

WHILE uiRxLength > 0 DO
    IF NOT M_ProcessReceive() THEN
        EXIT;
    END_IF
    bProcessed := TRUE;

    // 生成即时协议响应后,先交还主状态机发送 ACK,再继续处理后续入站帧。
    IF xPendingImmediateTx THEN
        EXIT;
    END_IF
END_WHILE

M_ProcessPendingFrames := bProcessed;

系列导航

  • 系列定位:第 4 篇
  • 上一篇:第3篇 PUBLISH 报文
  • 下一篇:第5篇 SUBSCRIBE / UNSUBSCRIBE 和主题过滤器校验
评论和回复区

评论区预留

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

↑ ↓