这一篇站在架构层看 FB_MqttClient,重点讲主状态机为什么才是系统骨架、iConnected 为什么是调度中心、接收链和发送链为什么必须拆开,以及异常为什么统一收口到 iTcpDisconnect。
适合谁收藏
- 正在做 CODESYS / PLC / MQTT 项目的人
- 想把 MQTT 从报文真正看到 ST 代码的人
- 正在排查 QoS1 / QoS2 超时、掉线、重连问题的人
如果你已经把前几篇都看完了,大概率会有一种感觉:
MQTT 报文其实不算太难,难的是这些报文怎么在 PLC 里有秩序地跑起来。
这个感觉是对的。
对 PLC 来说,真正决定一个 MQTT 客户端是“能跑”还是“跑稳”的,往往不是某个字节写错没写错,而是:
- 状态怎么分
- 状态怎么跳
- ACK 谁来发
- 超时谁来判
- 接收和发送怎么不互相卡死
这一篇我们就专门讲这件事。
先给结论:
报文是骨头,状态机才是筋。
一、先看 FB_MqttClient 的主状态全貌
直接上总图。
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 不是静止状态,而是 运行调度中心。
它至少要处理下面这些事:
- 有没有新的
Publish - 有没有新的
Subscribe - 有没有新的
Unsubscribe - 有没有待发送的即时 ACK
- 接收缓冲区里还有没有完整帧
- inflight 有没有超时项
- KeepAlive 有没有到期
- 用户是不是主动拉低
bConnect
所以更准确地说:
iConnected 是 MQTT 客户端的主循环枢纽。
三、为什么不能把所有收发逻辑都塞进一个状态里
有些 demo 风格实现,喜欢把事情写成这样:
- 连上后就在一个大状态里处理全部逻辑
- 收到了什么就现场回
- 要发什么就直接发
这种写法短期看起来省事,长期非常容易出问题。
为什么?
因为 MQTT 至少有 3 类动作在同时发生:
- 主动动作
- CONNECT
- PUBLISH
- SUBSCRIBE
- UNSUBSCRIBE
- PINGREQ
- 被动动作
- 对端发 PUBLISH 过来,你要回
PUBACK/PUBREC - 对端发
PUBREL过来,你要回PUBCOMP
- 后台动作
- inflight 超时扫描
- KeepAlive 计时
- 自动重连
如果这三类动作全堆在一个无边界的大状态里,后果通常就是:
- 某个 ACK 被拖延
- 某个超时条件没及时触发
- 某个等待状态迟迟不释放
四、为什么接收链和发送链必须分开
这是 MQTT 客户端设计里非常关键的一点。
发送链负责什么
- 主动构包
- 主动发出去
- 建立“我接下来在等谁”的等待关系
接收链负责什么
- 从接收缓冲区解析完整 MQTT 帧
- 判断报文类型
- 推进等待状态
- 必要时生成即时协议响应
把这两层拆开之后,你的代码脑子会清楚很多:
发送链决定“我要做什么”。
这个库里,这个分层就很明确:
| 职责 | 主要位置 |
|---|---|
| 主状态机总控 | FB_MqttClient.st |
| 报文发送构造 | 处理发送报文 |
| 报文接收解析 | M_ProcessReceive |
| 接收缓冲循环消化 | M_ProcessPendingFrames |
五、M_ProcessPendingFrames 为什么很关键
这个方法代码不长,但架构价值很高。
WHILE uiRxLength > 0 DO
IF NOT M_ProcessReceive() THEN
EXIT;
END_IF
bProcessed := TRUE;
IF xPendingImmediateTx THEN
EXIT;
END_IF
END_WHILE这段逻辑体现了 2 个重要思路:
- 只要缓冲区里还有完整帧,就持续消化
- 如果在处理中途生成了必须立刻回的 ACK,就先停下来,让主状态机优先把 ACK 发出去
这正是前面 QoS1 / QoS2 高频稳定性的关键之一。
如果你这里不分层、不分优先级,就容易出现:
- 入站帧堆在缓冲区
- ACK 迟迟发不出去
- 对端认为你超时了
六、为什么 iTcpDisconnect 是统一异常收口
一个成熟状态机,不会让每个异常都各自散落结束。 它通常会有一个统一异常收口状态。
在这个客户端里,就是:
iTcpDisconnect所有严重问题最终都可能导到这里,比如:
- TCP 连接超时
- CONNECT / PUBLISH / SUBSCRIBE 发送失败
- CONNACK / PUBACK / SUBACK 等待超时
- 即时 ACK 发送失败
这类设计的好处是:
- 断线清理逻辑集中
- 自动重连入口统一
- 诊断路径更清楚
这也是为什么现场你经常看到“最后都掉到 iTcpDisconnect”,但根因其实各不相同。
七、为什么高频场景主要在压 iConnected
前面我们一直在说高频 QoS1 / QoS2 最容易暴露问题。 如果从架构层面看,本质上就是在压 iConnected。
因为在这个状态里,客户端要同时协调:
- 收
- 发
- 等 ACK
- 发即时 ACK
- 扫 inflight
- 处理心跳
只要调度策略不清晰,这里就会先出问题。
所以判断一个 MQTT 客户端是否成熟,有个很实用的角度:
看它的 iConnected 是不是一个真正可调度、可回收、可让出优先级的运行态。
八、为什么说报文只是局部正确,状态机才决定整体正确
举个很实际的例子:
场景 1
你 PUBLISH 构包完全正确。 但 PUBACK 来了之后,状态机没有及时释放 xWaitingForAck。
结果:
- 报文没错
- 连接还是可能超时
场景 2
你 SUBSCRIBE 构包也完全正确。 但 SUBACK 收到后,状态没切回 iConnected。
结果:
- 报文没错
- 功能块外观看起来“卡死”
这就是为什么:
局部方法正确,只能说明“这块没写歪”。
九、读这套源码,最推荐的顺序是什么
如果你要真把 FB_MqttClient 看懂,我最建议按这个顺序读:
第一步:先看主状态骨架
FB_MqttClient.st
先建立“全局有哪些状态,状态怎么跳”的大图。
第二步:看 3 条主链
- CONNECT / CONNACK
- PUBLISH / QoS1 / QoS2
- SUBSCRIBE / UNSUBSCRIBE
第三步:回头看辅助机制
M_ProcessPendingFramesM_InflightCheckTimeoutM_IsValidTopicFilter
第四步:最后再看方法细节
否则一上来扎进方法内部,很容易只见树木,不见森林。
十、把主状态机的职责再压成一张表
| 状态段 | 核心职责 |
|---|---|
iDisconnected -> iConnAck | 建立 MQTT 会话 |
iConnected | 运行态调度中心 |
iPublish -> iPubComp | 发布与 QoS 确认链 |
iSubscribe -> iSubAck | 订阅链 |
iUnsubscribe -> iUnsubAck | 退订链 |
iPingReq -> iPingResp | 长连接保活 |
iTcpDisconnect | 异常收口与恢复入口 |
这张表对排障特别有用。 因为它让你一眼就知道:当前问题大概率属于哪一条链。
十一、这一篇你最该记住的 5 句话
- 报文决定局部动作,状态机决定整体行为。
iConnected不是空闲态,而是客户端运行调度中心。- 发送链和接收链必须分开,否则高频场景很容易互相阻塞。
M_ProcessPendingFrames这种小方法,往往决定客户端能不能把 ACK 时序跑顺。- 成熟 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 三个子方法调用。你会看到所有报文方法最终都是被主状态机统一调度的。
/// =======================================================================
/// 名称 : 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 现场排障:为什么会超时、掉线、重连、订阅丢失
评论区预留
这里先保留评论和回复结构,不接入第三方服务。后续统一决定登录、匿名、审核、反垃圾和静态站兼容策略。