这一篇专门讲高可靠链路为什么最容易出问题,核心内容包括 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 timeoutPubAck timeoutPubRec timeoutPubComp timeout- 状态机跳到
iTcpDisconnect
二、先把 QoS1 链路压成一张图
sequenceDiagram
participant PLC as PLC Client
participant Broker as Broker
PLC->>Broker: PUBLISH(QoS1)
Broker-->>PLC: PUBACK表面看只有两步。 但客户端内部实际至少要做下面这些事:
- 构造 PUBLISH
- 分配 Packet ID
- 建立 inflight
- 发送报文
- 标记“我正在等 PUBACK”
- 收到 PUBACK 后匹配 Packet ID
- 删除 inflight
- 释放等待状态
- 回到
iConnected
所以别被“图很短”骗了。 短的是网络报文,不短的是客户端内部控制逻辑。
三、QoS2 链路为什么更考验实现水平
QoS2 完整链路是:
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 放回状态机里看
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 做的。
真实代码非常值得看:
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 件事:
- 它不是盲目全量重发,而是扫描首个超时项
- 超时后会把
bDup置位 - 会累加重试次数
- 超过上限后,不再死循环,而是报错并清理
这就是为什么“能不能超时重发”不是一句空话,而是具体状态管理能力。
七、DUP 位到底什么时候才有意义
很多人知道 PUBLISH 首字节里有个 DUP 位。 但不知道它什么时候该置位。
结论很简单:
第一次发,不该乱置 DUP。
这个库里就是靠 inflight 超时后改 bDup 来驱动的。
也就是说,DUP 不是“高 QoS 就自动为 1”。 它表达的是:
这条消息不是第一次发,而是我重新发的。
八、为什么高频场景更容易暴露问题
高频时,客户端同时在干的事更多:
- 处理新的发布请求
- 等待上一条消息的 ACK
- 处理入站 ACK 报文
- 可能还要立即回对端 ACK
- 还要扫 inflight 超时
- KeepAlive 也可能同时到期
所以高频场景本质上是在压测两件事:
- 状态机是不是清晰
- 接收面和发送面有没有互相堵住
如果这两件事任何一个处理不好,就很容易出现一种现象:
明明网络没断、消息偶尔也收到了,但状态机却直接走到了 iTcpDisconnect。
这通常不是“网炸了”,而是协议时序没跟上。
九、为什么接收面即时消化这么关键
这个库有个很关键的方法:M_ProcessPendingFrames。
它的核心逻辑很短:
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. 状态机层
iPublishiPubAckiPubReciPubReliPubComp
4. 实现层
| 能力 | 方法 / 机制 |
|---|---|
| 构包 | M_BuildPublishPacket |
| 超时扫描 | M_InflightCheckTimeout |
| 接收帧循环消化 | M_ProcessPendingFrames |
| 真正解析 ACK / PUBLISH | M_ProcessReceive |
十三、这一篇你最该记住的 6 句话
- QoS1 / QoS2 的难点从来不是“发包”,而是“发出去以后怎么把后面的链走完”。
- 高可靠 MQTT 客户端必须维护 inflight,而不是只靠当前输入参数。
- 超时重发不只是再发一次,而是要带着状态、次数、DUP 位一起推进。
- 高频场景测的不是网速,测的是状态机调度能力。
- 很多 Immediate response timeout,不是没收到消息,而是没及时把协议 ACK 发回去。
- 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的影响。
/// =======================================================================
/// 名称 : 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 - 复制使用说明:这个方法决定哪一条在途消息超时、是否该触发重发,是高频场景稳定性的关键检查点。
- 阅读重点:重点看在途队列扫描、超时判定条件和返回索引,这一层直接决定状态机会不会进入重发分支。
/// =======================================================================
/// 名称 : 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 高一点就积帧、超时、掉线。
- 阅读重点:先看完整帧边界怎么判,再看循环消费逻辑,最后看为什么它能把“偶发来不及响应”变成“连续稳定处理”。
/// =======================================================================
/// 名称 : 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 和主题过滤器校验
评论区预留
这里先保留评论和回复结构,不接入第三方服务。后续统一决定登录、匿名、审核、反垃圾和静态站兼容策略。