ControlRookie
返回文章

第6篇_PLC 里写 MQTT,最难的不是报文,而是状态机

这一篇站在架构层看 FB_MqttClient,重点讲主状态机为什么才是系统骨架、iConnected 为什么是调度中心、接收链和发送链为什么必须拆开,以及异常为什么统一收口到 iTcpDisconnect。

这一篇站在架构层看 FB_MqttClient,重点讲主状态机为什么才是系统骨架、iConnected 为什么是调度中心、接收链和发送链为什么必须拆开,以及异常为什么统一收口到 iTcpDisconnect。

适合谁收藏

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

如果你已经把前几篇都看完了,大概率会有一种感觉:

MQTT 报文其实不算太难,难的是这些报文怎么在 PLC 里有秩序地跑起来。

这个感觉是对的。

对 PLC 来说,真正决定一个 MQTT 客户端是“能跑”还是“跑稳”的,往往不是某个字节写错没写错,而是:

  • 状态怎么分
  • 状态怎么跳
  • ACK 谁来发
  • 超时谁来判
  • 接收和发送怎么不互相卡死

这一篇我们就专门讲这件事。

先给结论:

报文是骨头,状态机才是筋。

一、先看 FB_MqttClient 的主状态全貌

直接上总图。

Mermaid
flowchart TD
    A[iDisconnected] --> B[iTcpConnect]
    B --> C[iConnect]
    C --> D[iConnAck]
    D --> E[iConnected]

    E --> F[iPublish]
    F --> G[iPubAck]
    F --> H[iPubRec]
    H --> I[iPubRel]
    I --> J[iPubComp]

    E --> K[iSubscribe]
    K --> L[iSubAck]

    E --> M[iUnsubscribe]
    M --> N[iUnsubAck]

    E --> O[iPingReq]
    O --> P[iPingResp]

    E --> Q[iDisconnect]
    B --> R[iTcpDisconnect]
    C --> R
    D --> R
    F --> R
    G --> R
    H --> R
    I --> R
    J --> R
    L --> R
    N --> R
    P --> R
    Q --> R

这张图先记住一个原则:

这个状态机不是按“代码文件夹”分的,而是按“协议动作链”分的。

也就是说:

  • 连接是一条链
  • 发布是一条链
  • 订阅是一条链
  • 心跳是一条链
  • 异常恢复是一条统一收口链

二、为什么 iConnected 是整个系统的调度中心

很多人第一次看状态机,会误以为 iConnected 是个“空闲态”。

其实恰恰相反。

在这个客户端里,iConnected 不是静止状态,而是 运行调度中心。

它至少要处理下面这些事:

  1. 有没有新的 Publish
  2. 有没有新的 Subscribe
  3. 有没有新的 Unsubscribe
  4. 有没有待发送的即时 ACK
  5. 接收缓冲区里还有没有完整帧
  6. inflight 有没有超时项
  7. KeepAlive 有没有到期
  8. 用户是不是主动拉低 bConnect

所以更准确地说:

iConnected 是 MQTT 客户端的主循环枢纽。

三、为什么不能把所有收发逻辑都塞进一个状态里

有些 demo 风格实现,喜欢把事情写成这样:

  • 连上后就在一个大状态里处理全部逻辑
  • 收到了什么就现场回
  • 要发什么就直接发

这种写法短期看起来省事,长期非常容易出问题。

为什么?

因为 MQTT 至少有 3 类动作在同时发生:

  1. 主动动作
  • CONNECT
  • PUBLISH
  • SUBSCRIBE
  • UNSUBSCRIBE
  • PINGREQ
  1. 被动动作
  • 对端发 PUBLISH 过来,你要回 PUBACK/PUBREC
  • 对端发 PUBREL 过来,你要回 PUBCOMP
  1. 后台动作
  • inflight 超时扫描
  • KeepAlive 计时
  • 自动重连

如果这三类动作全堆在一个无边界的大状态里,后果通常就是:

  • 某个 ACK 被拖延
  • 某个超时条件没及时触发
  • 某个等待状态迟迟不释放

四、为什么接收链和发送链必须分开

这是 MQTT 客户端设计里非常关键的一点。

发送链负责什么

  • 主动构包
  • 主动发出去
  • 建立“我接下来在等谁”的等待关系

接收链负责什么

  • 从接收缓冲区解析完整 MQTT 帧
  • 判断报文类型
  • 推进等待状态
  • 必要时生成即时协议响应

把这两层拆开之后,你的代码脑子会清楚很多:

发送链决定“我要做什么”。

这个库里,这个分层就很明确:

职责主要位置
主状态机总控FB_MqttClient.st
报文发送构造处理发送报文
报文接收解析M_ProcessReceive
接收缓冲循环消化M_ProcessPendingFrames

五、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

这段逻辑体现了 2 个重要思路:

  1. 只要缓冲区里还有完整帧,就持续消化
  2. 如果在处理中途生成了必须立刻回的 ACK,就先停下来,让主状态机优先把 ACK 发出去

这正是前面 QoS1 / QoS2 高频稳定性的关键之一。

如果你这里不分层、不分优先级,就容易出现:

  • 入站帧堆在缓冲区
  • ACK 迟迟发不出去
  • 对端认为你超时了

六、为什么 iTcpDisconnect 是统一异常收口

一个成熟状态机,不会让每个异常都各自散落结束。 它通常会有一个统一异常收口状态。

在这个客户端里,就是:

text
iTcpDisconnect

所有严重问题最终都可能导到这里,比如:

  • TCP 连接超时
  • CONNECT / PUBLISH / SUBSCRIBE 发送失败
  • CONNACK / PUBACK / SUBACK 等待超时
  • 即时 ACK 发送失败

这类设计的好处是:

  1. 断线清理逻辑集中
  2. 自动重连入口统一
  3. 诊断路径更清楚

这也是为什么现场你经常看到“最后都掉到 iTcpDisconnect”,但根因其实各不相同。


七、为什么高频场景主要在压 iConnected

前面我们一直在说高频 QoS1 / QoS2 最容易暴露问题。 如果从架构层面看,本质上就是在压 iConnected。

因为在这个状态里,客户端要同时协调:

  • 收
  • 发
  • 等 ACK
  • 发即时 ACK
  • 扫 inflight
  • 处理心跳

只要调度策略不清晰,这里就会先出问题。

所以判断一个 MQTT 客户端是否成熟,有个很实用的角度:

看它的 iConnected 是不是一个真正可调度、可回收、可让出优先级的运行态。

八、为什么说报文只是局部正确,状态机才决定整体正确

举个很实际的例子:

场景 1

你 PUBLISH 构包完全正确。 但 PUBACK 来了之后,状态机没有及时释放 xWaitingForAck。

结果:

  • 报文没错
  • 连接还是可能超时

场景 2

你 SUBSCRIBE 构包也完全正确。 但 SUBACK 收到后,状态没切回 iConnected。

结果:

  • 报文没错
  • 功能块外观看起来“卡死”

这就是为什么:

局部方法正确,只能说明“这块没写歪”。

九、读这套源码,最推荐的顺序是什么

如果你要真把 FB_MqttClient 看懂,我最建议按这个顺序读:

第一步:先看主状态骨架

  • FB_MqttClient.st

先建立“全局有哪些状态,状态怎么跳”的大图。

第二步:看 3 条主链

  1. CONNECT / CONNACK
  2. PUBLISH / QoS1 / QoS2
  3. SUBSCRIBE / UNSUBSCRIBE

第三步:回头看辅助机制

  • M_ProcessPendingFrames
  • M_InflightCheckTimeout
  • M_IsValidTopicFilter

第四步:最后再看方法细节

否则一上来扎进方法内部,很容易只见树木,不见森林。


十、把主状态机的职责再压成一张表

状态段核心职责
iDisconnected -> iConnAck建立 MQTT 会话
iConnected运行态调度中心
iPublish -> iPubComp发布与 QoS 确认链
iSubscribe -> iSubAck订阅链
iUnsubscribe -> iUnsubAck退订链
iPingReq -> iPingResp长连接保活
iTcpDisconnect异常收口与恢复入口

这张表对排障特别有用。 因为它让你一眼就知道:当前问题大概率属于哪一条链。


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

  1. 报文决定局部动作,状态机决定整体行为。
  2. iConnected 不是空闲态,而是客户端运行调度中心。
  3. 发送链和接收链必须分开,否则高频场景很容易互相阻塞。
  4. M_ProcessPendingFrames 这种小方法,往往决定客户端能不能把 ACK 时序跑顺。
  5. 成熟 MQTT 客户端的标志,不是状态多,而是状态边界清楚、异常收口统一。

十二、下篇预告

下一篇我们直接从工程现场最关心的视角切:

为什么会超时、掉线、重连、订阅丢失,到底该怎么查

也就是把前面所有报文、ACK、状态机知识,真正落到排障路径上。

这篇会非常实战。


完整 ST 代码

复制使用说明

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

代码阅读重点

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

完整代码 1:FB_MqttClient

  • 对应源码路径:10 MQTT/MqttClient_V1_0/Device/Application/MQTT/POUs/MqttClient NBS/FB_MqttClient.st
  • 复制使用说明:这一篇的核心就是整个主功能块本体。只看零碎方法理解不了状态机,必须把整个 FB_MqttClient 连起来看。
  • 阅读重点:优先看 VAR_INPUT/VAR_OUTPUT、CASE eState OF、超时定时器、TCP 三个子方法调用。你会看到所有报文方法最终都是被主状态机统一调度的。
iecst
/// =======================================================================
/// 名称      : FB_MqttClient
/// 功能      : MQTT 客户端(支持 MQTT 3.1.1 / MQTT 5.0)
/// 说明      : 实现 MQTT 连接、发布、订阅、接收、心跳与重连状态机
/// 编程人员  : ControlRookie
/// 时间      : 2026-01-10
/// 版本      : V2.1
/// =======================================================================
{attribute 'hide_all_locals'}
FUNCTION_BLOCK FB_MqttClient
VAR_INPUT
    // 连接配置
    bEnable                : BOOL := TRUE;                               // 使能客户端
    bConnect               : BOOL;                                       // 连接命令
    sBrokerIP              : STRING := '192.168.20.222';                // Broker IP 地址
    uiPort                 : UINT := 1883;                               // Broker 端口号
    sClientID              : STRING;                                     // 客户端标识符
    sUsername              : STRING;                                     // 用户名
    sPassword              : STRING;                                     // 密码

    // 协议版本
    eVersion               : E_MqttVersion := E_MqttVersion.byMqttVersion311; // MQTT 协议版本

    // 连接参数
    bCleanSession          : BOOL := TRUE;                               // 是否清理会话
    uiKeepAlive            : UINT := 60;                                 // 心跳周期(秒)
    bUseSSL                : BOOL := FALSE;                              // 是否启用 SSL
    bAutoReconnect         : BOOL := TRUE;                               // 是否自动重连
    uiReconnectDelay       : UINT := 5000;                               // 重连延时(毫秒)

    // MQTT 5.0 连接属性
    udiSessionExpiry       : UDINT := 0;                                 // 会话过期间隔
    uiReceiveMax           : UINT := 65535;                              // 最大接收数量
    udMaxPacketSize        : UDINT := 4096;                              // 最大报文长度
    bRequestResponseInfo   : BOOL := FALSE;                              // 是否请求响应信息
    bRequestProblemInfo    : BOOL := TRUE;                               // 是否请求问题信息

    // 遗嘱消息
    bWillFlag              : BOOL := FALSE;                              // 是否启用遗嘱消息
    bWillRetain            : BOOL := FALSE;                              // 遗嘱消息保留标志
    eWillQoS               : E_MqttQoS := E_MqttQoS.byQoS0;             // 遗嘱消息 QoS
    sWillTopic             : STRING;                                     // 遗嘱主题
    sWillMessage           : STRING;                                     // 遗嘱消息内容

    // 发布参数
    bPublish               : BOOL := FALSE;                              // 发布命令
    sPubTopic              : STRING := 'CodeSys';                        // 发布主题
    sPubPayload            : STRING := 'This is CodeSys';                // 发布载荷
    ePubQoS                : E_MqttQoS;                                  // 发布 QoS
    bPubRetain             : BOOL;                                       // 发布保留标志

    // 订阅参数
    bSubscribe             : BOOL := FALSE;                              // 订阅命令
    sSubTopic              : STRING := 'CodeSys';                        // 订阅主题
    eSubQoS                : E_MqttQoS;                                  // 订阅请求 QoS
    udiSubscriptionId      : UDINT;                                      // MQTT 5.0 订阅标识符(>0 时包含)
    bUnsubscribe           : BOOL;                                       // 取消订阅命令
    sUnsubTopic            : STRING := 'CodeSys';                        // 取消订阅主题
END_VAR
VAR_OUTPUT
    // 状态信息
    eState                 : E_MqttState;                                // 客户端当前状态
    bIsConnected           : BOOL;                                        // TCP 连接状态
    bMqttConnected         : BOOL;                                        // MQTT 协议连接状态
    bError                 : BOOL;                                        // 错误标志
    eErrorID               : NBS.ERROR;                                   // 错误码
    sDiagMsg               : STRING;                                      // 诊断信息

    // 订阅列表管理
    aSubscriptions         : ARRAY[1..GVL_Mqtt.cnMaxSubscriptions] OF ST_MqttSubscription; // 订阅列表
    uiSubscriptionCount    : UINT := 0;                                  // 当前订阅数量

    // 接收消息
    sRecTopic              : STRING;                                      // 最新接收主题
    sRecPayload            : STRING;                                      // 最新接收载荷
    aRecTopicList          : ARRAY[0..GVL_Mqtt.cnMaxHistory] OF STRING;   // 接收主题历史
    aRecPayloadList        : ARRAY[0..GVL_Mqtt.cnMaxHistory] OF STRING;   // 接收载荷历史
    byReceivedQoS          : BYTE;                                        // 最新接收消息 QoS
    bReceivedRetain        : BOOL;                                        // 最新接收消息保留标志
END_VAR
VAR
    rtrigConnect          : R_TRIG;                                      // 连接命令上升沿检测
    ftrigConnect          : F_TRIG;                                      // 连接命令下降沿检测
    rtrigPublish          : R_TRIG;                                      // 发布命令上升沿检测
    rtrigSubscribe        : R_TRIG;                                      // 订阅命令上升沿检测
    rtrigUnsubscribe      : R_TRIG;                                      // 取消订阅命令上升沿检测

    // 超时计时
    tTimeout                : TIME;                                       // 当前超时时间设定
    tonTimer                : TON;                                        // 通用超时定时器
    tonKeepAlive            : TON;                                        // 心跳定时器
    tonReconnect            : TON;                                        // 重连延时定时器
    eLastState              : E_MqttState;                                // 上一周期状态

    // TCP连接对象
    fbTcpClient              : NBS.TCP_Client;                             // TCP 客户端实例
    fbTcpRead                : NBS.TCP_Read;                               // TCP 读取实例
    fbTcpWrite               : NBS.TCP_Write;                              // TCP 写入实例
    hConnection              : NBS.CAA.HANDLE;                             // TCP 连接句柄
    bTcpConnect              : BOOL;                                       // TCP 连接使能标志
    bTcpRead                 : BOOL;                                       // TCP 读取使能标志
    bTcpWrite                : BOOL;                                       // TCP 写入使能标志
    bWriteDoneLatched        : BOOL;                                       // TCP 写入完成锁存标志
    bHasRead                 : BOOL;                                       // 本周期已读取完成标志
    bHasWritten              : BOOL;                                       // 本周期已写入完成标志
    udiBytesRead             : UDINT;                                      // 本次 TCP 读取字节数
    aTxBuf                   : ARRAY[0..GVL_Mqtt.cnSendBufferSize - 1] OF BYTE; // 发送缓冲区
    aRxBuf                   : ARRAY[0..GVL_Mqtt.cnRecvBufferSize - 1] OF BYTE; // 接收缓冲区
    uiTxLength               : UINT;                                       // 待发送字节数
    uiRxLength               : UINT;                                       // 已接收字节数

    // 系统时间
    stTimeZone               : Util.TimeZone := (iBias := 480);            // 时区设置
    uliSysTime               : ULINT;                                      // 当前系统时间

    // MQTT协议相关
    bDup                     : BOOL;                                       // DUP 重发标志
    uiPacketId               : UINT := 1;                                 // 下一个报文标识符
    uiQoS2PacketId           : UINT;                                       // QoS2 流程报文标识符
    uiExpectedPacketId       : UINT;                                       // 当前期待确认的报文标识符
    byExpectedMsgType        : BYTE;                                       // 当前期待的报文类型
    xWaitingForAck           : BOOL;                                       // 等待通用 ACK 标志
    xWaitingForSubAck        : BOOL;                                       // 等待 SUBACK 标志
    xWaitingForUnsubAck      : BOOL;                                       // 等待 UNSUBACK 标志
    bPingPending             : BOOL;                                       // 心跳响应等待标志
    uiSendQuota              : UINT := GVL_Mqtt.cnDefaultReceiveMax;       // MQTT 5.0 发送配额
    uiInflightCount          : UINT;                                       // 在途消息数量
    uiTopicAliasCount        : UINT;                                       // 主题别名数量
    uiNextTopicAlias         : UINT := 1;                                  // 下一个发送主题别名
    uiPendingSubPacketId     : UINT;                                       // 待确认订阅报文标识符
    uiPendingUnsubPacketId   : UINT;                                       // 待确认取消订阅报文标识符
    uiRetryInflightIndex     : UINT;                                       // 待重发在途消息索引
    uiRxInFlightQosCount     : UINT;                                       // 接收侧未完成 QoS>0 消息数量
    xPendingImmediateTx      : BOOL;                                       // 接收路径即时回包待发送标志
    aInflight                : ARRAY[1..GVL_Mqtt.cnMaxInflight] OF ST_MqttInflightMessage; // 出站在途队列
    aTopicAlias              : ARRAY[1..GVL_Mqtt.cnMaxTopicAlias] OF ST_MqttTopicAlias; // 主题别名表
    aRxQoS2PacketIds         : ARRAY[1..GVL_Mqtt.cnMaxInflight] OF UINT;   // 入站 QoS2 去重表

    // 事件标志
    xConnectedEvent          : BOOL;                                       // 连接成功事件
    xDisconnectedEvent       : BOOL;                                       // 断开连接事件
    xSubscribedEvent         : BOOL;                                       // 订阅成功事件
    xUnsubscribedEvent       : BOOL;                                       // 取消订阅成功事件
    xPublishedEvent          : BOOL;                                       // 发布成功事件
    bMessageReceived         : BOOL;                                       // 收到消息事件

    // 统计信息
    uiMessagesSent           : UDINT;                                      // 已发送消息数量
    uiMessagesReceived       : UDINT;                                      // 已接收消息数量
    dtLastMessageTime        : DATE_AND_TIME;                              // 最后消息时间

    // 重连管理
    uiReconnectAttempts      : UINT := 0;                                  // 已重连次数
    uiMaxReconnectAttempts   : UINT := 10;                                 // 最大重连次数

    // MQTT 5.0 服务器属性(CONNACK解析后存储)
    uiServerReceiveMax       : UINT := 65535;                              // 服务端最大接收数量
    byServerMaxQoS           : BYTE := 2;                                  // 服务端支持的最大 QoS
    bServerRetainAvailable   : BOOL := TRUE;                               // 服务端是否支持保留消息
    udServerMaxPacketSize    : UDINT := GVL_Mqtt.cnMaxPacketSize;          // 服务端允许的最大报文长度
    uiServerTopicAliasMax    : UINT := 0;                                  // 服务端允许的最大主题别名
    bServerWildcardSubAvail  : BOOL := TRUE;                               // 服务端是否支持通配符订阅
    bServerSubIdAvail        : BOOL := TRUE;                               // 服务端是否支持订阅标识符
    bServerSharedSubAvail    : BOOL := TRUE;                               // 服务端是否支持共享订阅

END_VAR

// === IMPLEMENTATION ===
IF NOT bEnable AND (eState <> E_MqttState.iDisconnected) THEN
    IF bMqttConnected AND bIsConnected THEN
        eState := E_MqttState.iDisconnect;
    END_IF

    IF NOT bMqttConnected AND bIsConnected THEN
        eState := E_MqttState.iTcpDisconnect;
    END_IF
    tTimeout := T#0S;

    xConnectedEvent := FALSE;
    xDisconnectedEvent := FALSE;
    xSubscribedEvent := FALSE;
    xPublishedEvent := FALSE;

    uiMessagesSent := 0;
    uiMessagesReceived := 0;
    uiPacketId := 0;
    uiQoS2PacketId := 0;
    uiExpectedPacketId := 0;
    uiPendingSubPacketId := 0;
    uiPendingUnsubPacketId := 0;
    byExpectedMsgType := 0;
    xWaitingForAck := FALSE;
    xWaitingForSubAck := FALSE;
    xWaitingForUnsubAck := FALSE;
    bPingPending := FALSE;
    uiRetryInflightIndex := 0;
    M_InflightClear();
    M_TopicAliasClear();
    M_SubListClear();
    THIS^.M_ResetError();

    RETURN;
END_IF

/// 系统时间
stTimeZone.iBias := 480;
uliSysTime := GetLocalDateTime(tzTimeZone := stTimeZone);

/// 边沿检测
rtrigConnect(CLK := bConnect);
ftrigConnect(CLK := bConnect);
rtrigPublish(CLK := bPublish);
rtrigSubscribe(CLK := bSubscribe);
rtrigUnsubscribe(CLK := bUnsubscribe);

/// 清除单次事件标志
xConnectedEvent := FALSE;
xDisconnectedEvent := FALSE;
xSubscribedEvent := FALSE;
xUnsubscribedEvent := FALSE;
xPublishedEvent := FALSE;
bMessageReceived := FALSE;

/// 核心状态机
CASE eState OF
    //=======================================================================
    // 禁用状态
    //=======================================================================
    E_MqttState.iDisconnected:
        bTcpRead := FALSE;
        bTcpWrite := FALSE;

        IF rtrigConnect.Q OR (bAutoReconnect AND (uiReconnectAttempts > 0)) THEN
            IF sBrokerIP <> '' AND sClientID <> '' THEN
                bError := FALSE;
                eErrorID := TO_INT(E_ReasonCode.uiErrNoError);
                sDiagMsg := '';
                eState := E_MqttState.iTcpConnect;
            ELSE
                M_SetError(TO_UINT(E_ReasonCode.uiErrInvalidParameter), 'Invalid IP or ClientId');
            END_IF
        END_IF

    //=======================================================================
    // TCP连接中
    //=======================================================================
    E_MqttState.iTcpConnect:
        bTcpConnect := TRUE;

        IF bIsConnected AND hConnection <> 0 THEN
            eState := E_MqttState.iConnect;
        ELSIF fbTcpClient.xError THEN
            M_SetError(TO_UINT(E_ReasonCode.uiErrTcpConnectFailed), CONCAT('TCP error: ', INT_TO_STRING(fbTcpClient.eError)));
            eState := E_MqttState.iTcpDisconnect;
        ELSE
            tTimeout := T#5S;
            IF tonTimer.Q THEN
                M_SetError(TO_UINT(E_ReasonCode.uiErrTimeout), 'TcpConnect timeout');
                eState := E_MqttState.iTcpDisconnect;
            END_IF
        END_IF

    //=======================================================================
    // 发送MQTT CONNECT报文
    //=======================================================================
    E_MqttState.iConnect:
        IF NOT bHasWritten THEN
            IF NOT M_BuildConnectPacket() THEN
                eState := E_MqttState.iTcpDisconnect;
            END_IF
        END_IF

        IF eState = E_MqttState.iConnect THEN
            bTcpWrite := TRUE;
        END_IF
        IF (eState = E_MqttState.iConnect) AND bHasWritten THEN
            bTcpWrite := FALSE;
            xWaitingForAck := TRUE;
            byExpectedMsgType := E_MqttPacketType.byConnAck;
            eState := E_MqttState.iConnAck;
        ELSIF (eState = E_MqttState.iConnect) AND fbTcpWrite.xError THEN
            M_SetError(TO_UINT(E_ReasonCode.uiErrTcpSendFailed), CONCAT('Send error: ', INT_TO_STRING(fbTcpWrite.eError)));
            eState := E_MqttState.iTcpDisconnect;
        ELSIF eState = E_MqttState.iConnect THEN
            tTimeout := T#2S;
            IF tonTimer.Q THEN
                M_SetError(TO_UINT(E_ReasonCode.uiErrTimeout), 'MQTT connection timeout');
                eState := E_MqttState.iTcpDisconnect;
            END_IF
        END_IF

    //=======================================================================
    // 等待CONNACK响应
    //=======================================================================
    E_MqttState.iConnAck:
        tTimeout := T#2S;

        IF (uiRxLength > 0) OR M_ReadIntoBuffer() THEN
            tTimeout := T#0S;
            IF M_HandleConnAck() THEN
                xConnectedEvent := TRUE;
                tonKeepAlive(IN := FALSE);
                uiReconnectAttempts := 0;
                eState := E_MqttState.iConnected;
            ELSE
                M_SetError(TO_UINT(E_ReasonCode.uiErrConnAckRefused), sDiagMsg);
                eState := E_MqttState.iTcpDisconnect;
            END_IF
        ELSIF fbTcpRead.xError THEN
            M_SetError(TO_UINT(E_ReasonCode.uiErrTcpReceiveFailed), CONCAT('Receive error: ', INT_TO_STRING(fbTcpRead.eError)));
            eState := E_MqttState.iTcpDisconnect;
        ELSIF tonTimer.Q THEN
            M_SetError(TO_UINT(E_ReasonCode.uiErrTimeout), 'ConnAck timeout');
            eState := E_MqttState.iTcpDisconnect;
        END_IF

    //=======================================================================
    // 已连接状态
    //=======================================================================
    E_MqttState.iConnected:
        IF NOT bIsConnected THEN
            M_SetError(TO_UINT(E_ReasonCode.uiErrTimeout), 'TCP disconnected');
            eState := E_MqttState.iTcpDisconnect;
        END_IF

        IF (eState = E_MqttState.iConnected) AND xPendingImmediateTx THEN
            bTcpWrite := TRUE;
            tTimeout := GVL_Mqtt.cnResponseTimeout;
            IF bHasWritten THEN
                bTcpWrite := FALSE;
                tTimeout := T#0S;
                IF uiRxInFlightQosCount > 0 THEN
                    IF (uiTxLength > 0) AND ((aTxBuf[0] AND GVL_Mqtt.cnHdrTypeMask) = E_MqttPacketType.byPubAck) THEN
                        uiRxInFlightQosCount := uiRxInFlightQosCount - 1;
                    END_IF
                END_IF
                uiTxLength := 0;
                xPendingImmediateTx := FALSE;
            ELSIF fbTcpWrite.xError THEN
                bTcpWrite := FALSE;
                tTimeout := T#0S;
                xPendingImmediateTx := FALSE;
                M_SetError(TO_UINT(E_ReasonCode.uiErrTcpSendFailed), 'Immediate response send failed');
                eState := E_MqttState.iTcpDisconnect;
            ELSIF tonTimer.Q THEN
                bTcpWrite := FALSE;
                tTimeout := T#0S;
                xPendingImmediateTx := FALSE;
                M_SetError(TO_UINT(E_ReasonCode.uiErrTimeout), 'Immediate response timeout');
                eState := E_MqttState.iTcpDisconnect;
            END_IF
        ELSIF eState = E_MqttState.iConnected THEN
            tTimeout := T#0S;
        END_IF

        IF (eState = E_MqttState.iConnected) AND (NOT xPendingImmediateTx) THEN
            IF (uiRxLength > 0) OR M_ReadIntoBuffer() THEN
                tonKeepAlive(IN := FALSE);
                M_ProcessPendingFrames();
            END_IF
        END_IF

        IF (eState = E_MqttState.iConnected) AND (NOT xPendingImmediateTx) THEN
            uiRetryInflightIndex := M_InflightCheckTimeout();
            IF uiRetryInflightIndex > 0 THEN
                CASE aInflight[uiRetryInflightIndex].eState OF
                    E_MqttInflightState.iPublishSent:
                        ePubQoS := aInflight[uiRetryInflightIndex].eQoS;
                        eState := E_MqttState.iPublish;
                    E_MqttInflightState.iPubRelSent:
                        uiQoS2PacketId := aInflight[uiRetryInflightIndex].uiPacketId;
                        eState := E_MqttState.iPubRel;
                ELSE
                        uiRetryInflightIndex := 0;
                END_CASE
            END_IF

            IF tonKeepAlive.Q THEN
                tonKeepAlive(IN := FALSE);
                eState := E_MqttState.iPingReq;
            END_IF

            IF (eState = E_MqttState.iConnected) AND rtrigPublish.Q AND NOT xWaitingForAck AND NOT xWaitingForSubAck AND NOT xWaitingForUnsubAck AND NOT bPingPending THEN
                tonKeepAlive(IN := FALSE);
                uiRetryInflightIndex := 0;
                eState := E_MqttState.iPublish;
            END_IF

            IF (eState = E_MqttState.iConnected) AND rtrigSubscribe.Q AND NOT xWaitingForAck AND NOT xWaitingForSubAck AND NOT xWaitingForUnsubAck AND NOT bPingPending THEN
                tonKeepAlive(IN := FALSE);
                eState := E_MqttState.iSubscribe;
            END_IF

            IF (eState = E_MqttState.iConnected) AND rtrigUnsubscribe.Q AND NOT xWaitingForAck AND NOT xWaitingForSubAck AND NOT xWaitingForUnsubAck AND NOT bPingPending THEN
                tonKeepAlive(IN := FALSE);
                eState := E_MqttState.iUnsubscribe;
            END_IF

            IF (eState = E_MqttState.iConnected) AND ftrigConnect.Q THEN
                eState := E_MqttState.iDisconnect;
            END_IF
        END_IF

    //=======================================================================
    // 心跳请求
    //=======================================================================
    E_MqttState.iPingReq:
        IF NOT bHasWritten THEN
            IF NOT M_BuildPingReqPacket() THEN
                eState := E_MqttState.iTcpDisconnect;
            END_IF
        END_IF

        IF eState = E_MqttState.iPingReq THEN
            bTcpWrite := TRUE;
        END_IF
        IF (eState = E_MqttState.iPingReq) AND bHasWritten THEN
            bTcpWrite := FALSE;
            bPingPending := TRUE;
            eState := E_MqttState.iPingResp;
        ELSIF (eState = E_MqttState.iPingReq) AND fbTcpWrite.xError THEN
            bTcpWrite := FALSE;
            M_SetError(TO_UINT(E_ReasonCode.uiErrTcpSendFailed), 'PingReq send failed');
            eState := E_MqttState.iTcpDisconnect;
        ELSIF eState = E_MqttState.iPingReq THEN
            tTimeout := T#2S;
            IF tonTimer.Q THEN
                M_SetError(TO_UINT(E_ReasonCode.uiErrTimeout), 'PingReq timeout');
                eState := E_MqttState.iTcpDisconnect;
            END_IF
        END_IF

    E_MqttState.iPingResp:
        tTimeout := T#2S;
        IF (uiRxLength > 0) OR M_ReadIntoBuffer() THEN
            IF M_ProcessPendingFrames() AND NOT bPingPending THEN
                tTimeout := T#0S;
                eState := E_MqttState.iConnected;
            END_IF
        ELSIF fbTcpRead.xError THEN
            M_SetError(TO_UINT(E_ReasonCode.uiErrTcpReceiveFailed), 'PingResp receive failed');
            eState := E_MqttState.iTcpDisconnect;
        ELSIF tonTimer.Q THEN
            bPingPending := FALSE;
            M_SetError(TO_UINT(E_ReasonCode.uiErrKeepAliveTimeout), 'PingResp timeout');
            eState := E_MqttState.iTcpDisconnect;
        END_IF

    //=======================================================================
    // 发布消息
    //=======================================================================
    E_MqttState.iPublish:
        IF NOT bHasWritten THEN
            IF NOT M_BuildPublishPacket() THEN
                IF NOT bError THEN
                    M_SetError(TO_UINT(E_ReasonCode.uiErrInvalidParameter), 'Build publish packet failed');
                END_IF
                uiRetryInflightIndex := 0;
                eState := E_MqttState.iConnected;
            END_IF
        END_IF

        IF eState = E_MqttState.iPublish THEN
            bTcpWrite := TRUE;
        END_IF
        IF (eState = E_MqttState.iPublish) AND bHasWritten THEN
            bTcpWrite := FALSE;
            uiMessagesSent := uiMessagesSent + 1;
            dtLastMessageTime := ULINT_TO_DT(uliSysTime / 1000);

            CASE ePubQoS OF
                E_MqttQoS.byQoS0:
                    xPublishedEvent := TRUE;
                    eState := E_MqttState.iConnected;
                    uiRetryInflightIndex := 0;

                E_MqttQoS.byQoS1:
                    xWaitingForAck := TRUE;
                    byExpectedMsgType := E_MqttPacketType.byPubAck;
                    eState := E_MqttState.iPubAck;

                E_MqttQoS.byQoS2:
                    xWaitingForAck := TRUE;
                    byExpectedMsgType := E_MqttPacketType.byPubRec;
                    eState := E_MqttState.iPubRec;
            ELSE
                M_SetError(TO_UINT(E_ReasonCode.uiErrInvalidParameter), 'Invalid QoS level');
                eState := E_MqttState.iConnected;
            END_CASE
        ELSIF (eState = E_MqttState.iPublish) AND fbTcpWrite.xError THEN
            bTcpWrite := FALSE;
            M_SetError(TO_UINT(E_ReasonCode.uiErrTcpSendFailed), 'Publish send failed');
            eState := E_MqttState.iTcpDisconnect;
        ELSIF eState = E_MqttState.iPublish THEN
            tTimeout := GVL_Mqtt.cnPublishTimeout;
            IF tonTimer.Q THEN
                M_SetError(TO_UINT(E_ReasonCode.uiErrTimeout), 'MQTT Publish timeout');
                eState := E_MqttState.iTcpDisconnect;
            END_IF
        END_IF

    //=======================================================================
    // 等待PUBACK (QoS 1)
    //=======================================================================
    E_MqttState.iPubAck:
        tTimeout := T#2S;

        IF (uiRxLength > 0) OR M_ReadIntoBuffer() THEN
            tTimeout := T#0S;
            IF M_ProcessPendingFrames() AND NOT xWaitingForAck THEN
                tTimeout := T#0S;
                eState := E_MqttState.iConnected;
            END_IF
        ELSIF fbTcpRead.xError THEN
            M_SetError(TO_UINT(E_ReasonCode.uiErrTcpReceiveFailed), CONCAT('Receive error: ', INT_TO_STRING(fbTcpRead.eError)));
            eState := E_MqttState.iTcpDisconnect;
        ELSIF tonTimer.Q THEN
            M_SetError(TO_UINT(E_ReasonCode.uiErrTimeout), 'PubAck timeout');
            eState := E_MqttState.iTcpDisconnect;
        END_IF

    //=======================================================================
    // 发布收到 (QoS 2 - 步骤1)
    //=======================================================================
    E_MqttState.iPubRec:
        tTimeout := T#2S;
        IF (uiRxLength > 0) OR M_ReadIntoBuffer() THEN
            tTimeout := T#0S;
            IF M_ProcessPendingFrames() AND NOT xWaitingForAck THEN
                tTimeout := T#0S;
                eState := E_MqttState.iPubRel;
            END_IF
        ELSIF fbTcpRead.xError THEN
            M_SetError(TO_UINT(E_ReasonCode.uiErrTcpReceiveFailed), CONCAT('Receive error: ', INT_TO_STRING(fbTcpRead.eError)));
            eState := E_MqttState.iTcpDisconnect;
        ELSIF tonTimer.Q Then
            M_SetError(TO_UINT(E_ReasonCode.uiErrTimeout), 'PubRec timeout');
            eState := E_MqttState.iTcpDisconnect;
        END_IF

    //=======================================================================
    // 发布释放 (QoS 2 - 步骤2)
    //=======================================================================
    E_MqttState.iPubRel:
        IF NOT bHasWritten THEN
            IF NOT M_BuildPubRelPacket() THEN
                eState := E_MqttState.iTcpDisconnect;
            END_IF
        END_IF

        IF eState = E_MqttState.iPubRel THEN
            bTcpWrite := TRUE;
        END_IF
        IF (eState = E_MqttState.iPubRel) AND bHasWritten THEN
            bTcpWrite := FALSE;
            xWaitingForAck := TRUE;
            byExpectedMsgType := E_MqttPacketType.byPubComp;
            M_InflightUpdateState(
                uiPacketId := uiQoS2PacketId,
                eNewState := E_MqttInflightState.iPubRelSent);
            eState := E_MqttState.iPubComp;
        ELSIF (eState = E_MqttState.iPubRel) AND fbTcpWrite.xError THEN
            bTcpWrite := FALSE;
            M_SetError(TO_UINT(E_ReasonCode.uiErrTcpSendFailed), 'PubRel send failed');
            eState := E_MqttState.iTcpDisconnect;
        ELSIF eState = E_MqttState.iPubRel THEN
            tTimeout := T#2S;
            IF tonTimer.Q THEN
                M_SetError(TO_UINT(E_ReasonCode.uiErrTimeout), 'MQTT PubRel timeout');
                eState := E_MqttState.iTcpDisconnect;
            END_IF
        END_IF

    //=======================================================================
    // QoS 2 消息发布完成 (QoS 2 - 步骤3)
    //=======================================================================
    E_MqttState.iPubComp:
        tTimeout := T#2S;
        IF (uiRxLength > 0) OR M_ReadIntoBuffer() THEN
            tTimeout := T#0S;
            IF M_ProcessPendingFrames() AND NOT xWaitingForAck THEN
                tTimeout := T#0S;
                eState := E_MqttState.iConnected;
            END_IF
        ELSIF fbTcpRead.xError THEN
            M_SetError(TO_UINT(E_ReasonCode.uiErrTcpReceiveFailed), CONCAT('Receive error: ', INT_TO_STRING(fbTcpRead.eError)));
            eState := E_MqttState.iTcpDisconnect;
        ELSIF tonTimer.Q THEN
            M_SetError(TO_UINT(E_ReasonCode.uiErrTimeout), 'PubComp timeout');
            eState := E_MqttState.iTcpDisconnect;
        END_IF

    //=======================================================================
    // 客户端订阅请求
    //=======================================================================
    E_MqttState.iSubscribe:
        IF NOT bHasWritten THEN
            IF NOT M_BuildSubscribePacket() THEN
                eState := E_MqttState.iConnected;
            END_IF
        END_IF

        IF eState = E_MqttState.iSubscribe THEN
            bTcpWrite := TRUE;
        END_IF
        IF (eState = E_MqttState.iSubscribe) AND bHasWritten THEN
            bTcpWrite := FALSE;
            xWaitingForSubAck := TRUE;
            byExpectedMsgType := E_MqttPacketType.bySubAck;
            eState := E_MqttState.iSubAck;

        ELSIF (eState = E_MqttState.iSubscribe) AND fbTcpWrite.xError THEN
            bTcpWrite := FALSE;
            M_SetError(TO_UINT(E_ReasonCode.uiErrTcpSendFailed), 'Subscribe send failed');
            eState := E_MqttState.iTcpDisconnect;
        ELSIF eState = E_MqttState.iSubscribe THEN
            tTimeout := T#2S;
            IF tonTimer.Q Then
                M_SetError(TO_UINT(E_ReasonCode.uiErrTimeout), 'MQTT Subscribe timeout');
                eState := E_MqttState.iTcpDisconnect;
            END_IF
        END_IF

    E_MqttState.iSubAck:
        tTimeout := T#2S;
        IF (uiRxLength > 0) OR M_ReadIntoBuffer() THEN
            IF M_ProcessPendingFrames() AND NOT xWaitingForSubAck THEN
                tTimeout := T#0S;
                eState := E_MqttState.iConnected;
            END_IF
        ELSIF fbTcpRead.xError THEN
            M_SetError(TO_UINT(E_ReasonCode.uiErrTcpReceiveFailed), 'SubAck receive failed');
            eState := E_MqttState.iTcpDisconnect;
        ELSIF tonTimer.Q THEN
            xWaitingForSubAck := FALSE;
            M_SetError(TO_UINT(E_ReasonCode.uiErrTimeout), 'SubAck timeout');
            eState := E_MqttState.iTcpDisconnect;
        END_IF

    //=======================================================================
    // 客户端取消订阅请求
    //=======================================================================
    E_MqttState.iUnsubscribe:
        IF NOT bHasWritten THEN
            IF NOT M_BuildUnsubscribePacket() THEN
                eState := E_MqttState.iConnected;
            END_IF
        END_IF

        IF eState = E_MqttState.iUnsubscribe THEN
            bTcpWrite := TRUE;
        END_IF
        IF (eState = E_MqttState.iUnsubscribe) AND bHasWritten Then
            bTcpWrite := FALSE;
            xWaitingForUnsubAck := TRUE;
            byExpectedMsgType := E_MqttPacketType.byUnsubAck;
            eState := E_MqttState.iUnsubAck;

        ELSIF (eState = E_MqttState.iUnsubscribe) AND fbTcpWrite.xError THEN
            bTcpWrite := FALSE;
            M_SetError(TO_UINT(E_ReasonCode.uiErrTcpSendFailed), 'Unsubscribe send failed');
            eState := E_MqttState.iTcpDisconnect;
        ELSIF eState = E_MqttState.iUnsubscribe THEN
            tTimeout := T#2S;
            IF tonTimer.Q THEN
                M_SetError(TO_UINT(E_ReasonCode.uiErrTimeout), 'MQTT Unsubscribe timeout');
                eState := E_MqttState.iTcpDisconnect;
            END_IF
        END_IF

    E_MqttState.iUnsubAck:
        tTimeout := T#2S;
        IF (uiRxLength > 0) OR M_ReadIntoBuffer() THEN
            IF M_ProcessPendingFrames() AND NOT xWaitingForUnsubAck THEN
                tTimeout := T#0S;
                eState := E_MqttState.iConnected;
            END_IF
        ELSIF fbTcpRead.xError THEN
            M_SetError(TO_UINT(E_ReasonCode.uiErrTcpReceiveFailed), 'UnsubAck receive failed');
            eState := E_MqttState.iTcpDisconnect;
        ELSIF tonTimer.Q THEN
            xWaitingForUnsubAck := FALSE;
            M_SetError(TO_UINT(E_ReasonCode.uiErrTimeout), 'UnsubAck timeout');
            eState := E_MqttState.iTcpDisconnect;
        END_IF

    //=======================================================================
    // 客户端断开连接
    //=======================================================================
    E_MqttState.iDisconnect:
        IF NOT bHasWritten THEN
            IF NOT M_BuildDisconnectPacket() THEN
                eState := E_MqttState.iTcpDisconnect;
            END_IF
        END_IF

        IF eState = E_MqttState.iDisconnect THEN
            bTcpWrite := TRUE;
        END_IF
        IF (eState = E_MqttState.iDisconnect) AND bHasWritten THEN
            bTcpWrite := FALSE;
            eState := E_MqttState.iTcpDisconnect;

        ELSIF (eState = E_MqttState.iDisconnect) AND fbTcpWrite.xError THEN
            bTcpWrite := FALSE;
            M_SetError(TO_UINT(E_ReasonCode.uiErrTcpSendFailed), 'Disconnect send failed');
            eState := E_MqttState.iTcpDisconnect;
        ELSIF eState = E_MqttState.iDisconnect THEN
            tTimeout := T#2S;
            IF tonTimer.Q THEN
                M_SetError(TO_UINT(E_ReasonCode.uiErrTimeout), 'MQTT Disconnect timeout');
                eState := E_MqttState.iTcpDisconnect;
            END_IF
        END_IF

    //=======================================================================
    // TCP断开连接
    //=======================================================================
    E_MqttState.iTcpDisconnect:
        bTcpConnect := FALSE;
        bTcpRead := FALSE;
        bTcpWrite := FALSE;
        xDisconnectedEvent := TRUE;
        xWaitingForAck := FALSE;
        xWaitingForSubAck := FALSE;
        xWaitingForUnsubAck := FALSE;
        bPingPending := FALSE;
        xPendingImmediateTx := FALSE;
        uiPendingSubPacketId := 0;
        uiPendingUnsubPacketId := 0;
        uiRetryInflightIndex := 0;

        // 自动重连逻辑
        IF bAutoReconnect AND uiReconnectAttempts < uiMaxReconnectAttempts THEN
            tonReconnect(IN := TRUE, PT := UINT_TO_TIME(uiReconnectDelay));
            IF tonReconnect.Q THEN
                tonReconnect(IN := FALSE);
                uiReconnectAttempts := uiReconnectAttempts + 1;
                eState := E_MqttState.iDisconnected;
            END_IF
        ELSE
            tonReconnect(IN := FALSE);
            eState := E_MqttState.iDisconnected;
        END_IF
ELSE
        M_SetError(TO_UINT(E_ReasonCode.uiErrInvalidState), 'Invalid state');
        eState := E_MqttState.iDisconnected;
END_CASE

/// =======================================================================
/// 标志位
/// =======================================================================
bMqttConnected S= xConnectedEvent;
bMqttConnected R= xDisconnectedEvent;

/// =======================================================================
/// 超时计时器
/// =======================================================================
tonTimer(
    IN := (eLastState = eState) AND (tTimeout <> T#0S),
    PT := tTimeout);
eLastState := eState;

/// KeepAlive心跳(BUG-09修复:在所有连接状态下运行)
tonKeepAlive(
    IN := (uiKeepAlive <> 0) AND bIsConnected
          AND (eState <> E_MqttState.iDisconnected)
          AND (eState <> E_MqttState.iTcpConnect)
          AND (eState <> E_MqttState.iTcpDisconnect),
    PT := UINT_TO_TIME(uiKeepAlive * 1000));

/// =======================================================================
/// TCP
/// =======================================================================
M_TcpClient(
    bEnable := bTcpConnect,
    sIP := sBrokerIP,
    uiPortNum := uiPort,
    bIsConnected => bIsConnected);

M_TcpRead(
    bEnable := bTcpRead,
    pDataReceive := ADR(aRxBuf[uiRxLength]),
    udiDataSize := SIZEOF(aRxBuf) - uiRxLength,
    bDone => bHasRead,
    udiBytesRead => udiBytesRead);

M_TcpWrite(
    bExecute := bTcpWrite,
    pDataSend := ADR(aTxBuf),
    udiDataSize := uiTxLength,
    bDone => bHasWritten);

系列导航

  • 系列定位:第 6 篇
  • 上一篇:第5篇 SUBSCRIBE / UNSUBSCRIBE
  • 下一篇:第7篇 MQTT 现场排障:为什么会超时、掉线、重连、订阅丢失
评论和回复区

评论区预留

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

↑ ↓