乐于分享
好东西不私藏

源码加更04_MQTT 编解码器和字节工具函数

源码加更04_MQTT 编解码器和字节工具函数

源码加更04_MQTT 编解码器和字节工具函数

这一组源码加更只有一个目标:把 MqttBroker 的真实 ST 源码按工程阅读顺序讲完整。不是再补几段“看起来像源码”的片段,而是让读者能沿着源码对象理解这个 Broker 怎么组织、怎么运行、怎么排障。

适合谁收藏

  • 已经读过 MqttBroker 主线教程,想继续看真实源码实现的工程师。
  • 想学习 CodeSys ST 工程如何拆分 Broker、连接池、编解码、路由和 QoS 调度的人。
  • 想把 MQTT Broker 移植到 PLC、边缘控制器或教学工程里的开发者。

先给结论

这一篇把 MQTT 报文解析和构造集中到 Codec 与工具函数里看,重点是从字节流还原协议字段,以及从结构化字段重新构造 ACK/PUBLISH。

MQTT 报文不是字符串拼接。Remaining Length、UTF-8 String、PacketId、QoS 标志位这些字节级边界,一旦错一个,Broker 就会表现成偶发断开或订阅失败。

这篇覆盖 16 个源码文件,合计约 1616 行 ST 代码。为了保持公开教程可读性,正文先讲源码阅读路径,再给完整源码。读代码时建议不要从第一个代码块一路机械读到底,而是按本篇的“读代码顺序”来抓主线。

从工程问题到代码职责

层次
本篇重点
你读源码时要抓住的判断
工程入口
程序如何启动、对象如何被实例化
先确认谁是入口,谁只是被调度的对象
数据边界
容量、状态、错误、缓冲区和表结构
先知道边界,后面排障才不会乱猜
协作关系
各 FB、函数和结构体如何互相传递数据
不按文件夹读,按数据流和状态流读
验证路径
在线观察应该看哪些变量
代码最终要能落到现场排障,而不是只停在源码阅读

本篇源码覆盖表

序号
源码对象
行数
1
FB_MqttBrokerCodec.M_BuildPublish.st
146
2
FB_MqttBrokerCodec.M_BuildSimpleAck.st
277
3
FB_MqttBrokerCodec.M_ParseConnect.st
290
4
FB_MqttBrokerCodec.M_ParsePublish.st
143
5
FB_MqttBrokerCodec.M_ParseSubscribe.st
145
6
FB_MqttBrokerCodec.M_ParseUnsubscribe.st
122
7
FB_MqttBrokerCodec.st
14
8
F_MqttAppendString.st
49
9
F_MqttContainsWildcard.st
34
10
F_MqttDecodeRemainingLength.st
63
11
F_MqttEncodeRemainingLength.st
55
12
F_MqttIsValidTopicFilter.st
76
13
F_MqttIsValidTopicName.st
39
14
F_MqttReadString.st
67
15
F_MqttSkipVariableByteInteger.st
54
16
F_MqttStartsWith.st
42

推荐阅读顺序

  • 先看 FB_MqttBrokerCodec.st 主体职责。
  • 再看 CONNECT/PUBLISH/SUBSCRIBE/UNSUBSCRIBE 解析。
  • 最后看 Remaining Length、String、Topic、QoS 等工具函数。

验证和排障边界

  • 订阅失败、报文断开、PacketId 不匹配时,优先检查本篇对象。
  • 抓包对照固定报头和 Remaining Length,可以最快定位编码边界错误。

本篇完整开源代码

下面代码来自对应 .st 源文件的连续完整内容。为方便公开阅读,只保留源码对象名,不放本机工程路径。

完整代码 01: FB_MqttBrokerCodec.M_BuildPublish.st

/// =======================================================================/// 名称      : M_BuildPublish/// 功能      : 构建 Broker 出站 PUBLISH 报文/// 说明      : 根据路由后的发布帧生成投递给订阅者的 MQTT PUBLISH。/// 编程人员  : ControlRookie/// 时间      : 2026-05-08/// 版本      : V1.0/// ======================================================================={attribute 'hide_all_locals'}METHOD M_BuildPublish : BOOLVAR_INPUT    stPublish     : ST_MqttBrokerPublishFrame; // 待投递给订阅者的发布帧    byProtocolLevel : BYTE; // 目标客户端 MQTT 协议级别,5 表示 PUBLISH 可变头需要追加零属性长度    uiWriteOffset : UINT; // 当前 PUBLISH 帧写入发送缓冲区的起始偏移,批量组包时用于追加到上一帧之后[byte]    udiBufferSize : UDINT; // 发送缓冲区总容量[byte]END_VARVAR_IN_OUT    aBuffer       : ARRAY[*] OF BYTE; // MQTT 发送缓冲区END_VARVAR_OUTPUT    uiFrameLen    : UINT; // 构建出的 PUBLISH 报文长度[byte]END_VARVAR    aRemaining    : ARRAY[0..3] OF BYTE; // Remaining Length 编码临时缓冲区    uiRemainingLen : UINT; // Remaining Length 编码字节数[byte]    udiRemaining  : UDINT; // PUBLISH 剩余长度数值[byte]    udiFrameLen   : UDINT; // 当前 PUBLISH 完整 MQTT 帧长度,用于偏移写入前的总边界检查[byte]    uiOffset      : UINT; // 当前写入偏移[byte]    uiIndex       : UINT; // 字节复制索引[byte]END_VAR// === IMPLEMENTATION ===uiFrameLen := 0;IF NOT stPublish.xValid THEN    M_BuildPublish := FALSE;    RETURN;END_IFIF NOT F_MqttIsValidTopicName(sTopic := stPublish.sTopic) THEN    M_BuildPublish := FALSE;    RETURN;END_IFudiRemaining := 2 + TO_UDINT(stPublish.uiTopicLen) + TO_UDINT(stPublish.uiPayloadLen);IF stPublish.eQoS <> E_MqttQoS.byQoS0 THEN    udiRemaining := udiRemaining + 2;END_IFIF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN    // MQTT 5.0 PUBLISH 可变头在 Topic/PacketId 后必须携带 Properties。    // 当前 Broker 不发送任何 5.0 属性,因此属性长度固定编码为 0,占 1 字节。    udiRemaining := udiRemaining + 1;END_IFIF NOT F_MqttEncodeRemainingLength(    aBuffer := aRemaining,    udiValue := udiRemaining,    udiBufferSize := SIZEOF(aRemaining),    uiEncodedLen => uiRemainingLen) THEN    M_BuildPublish := FALSE;    RETURN;END_IFudiFrameLen := 1 + TO_UDINT(uiRemainingLen) + udiRemaining;IF (TO_UDINT(uiWriteOffset) + udiFrameLen) > udiBufferSize THEN    M_BuildPublish := FALSE;    RETURN;END_IFaBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.byPublish);IF stPublish.xDup THEN    aBuffer[uiWriteOffset] := aBuffer[uiWriteOffset] OR 16#08;END_IFCASE stPublish.eQoS OF    E_MqttQoS.byQoS0:        // QoS0 固定头 QoS 位保持 00。    E_MqttQoS.byQoS1:        aBuffer[uiWriteOffset] := aBuffer[uiWriteOffset] OR 16#02;    E_MqttQoS.byQoS2:        aBuffer[uiWriteOffset] := aBuffer[uiWriteOffset] OR 16#04;ELSE    M_BuildPublish := FALSE;    RETURN;END_CASEIF stPublish.xRetain THEN    aBuffer[uiWriteOffset] := aBuffer[uiWriteOffset] OR 16#01;END_IFuiIndex := 0;WHILE uiIndex < uiRemainingLen DO    aBuffer[uiWriteOffset + 1 + uiIndex] := aRemaining[uiIndex];    uiIndex := uiIndex + 1;END_WHILEuiOffset := uiWriteOffset + 1 + uiRemainingLen;IF NOT F_MqttAppendString(    aBuffer := aBuffer,    uiOffset := uiOffset,    sValue := stPublish.sTopic,    udiBufferSize := udiBufferSize) THEN    M_BuildPublish := FALSE;    RETURN;END_IFIF stPublish.eQoS <> E_MqttQoS.byQoS0 THEN    IF (TO_UDINT(uiOffset) + 2) > udiBufferSize THEN        M_BuildPublish := FALSE;        RETURN;    END_IF    aBuffer[uiOffset] := TO_BYTE(stPublish.uiPacketId / 256);    aBuffer[uiOffset + 1] := TO_BYTE(stPublish.uiPacketId MOD 256);    uiOffset := uiOffset + 2;END_IFIF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN    IF (TO_UDINT(uiOffset) + 1) > udiBufferSize THEN        M_BuildPublish := FALSE;        RETURN;    END_IF    aBuffer[uiOffset] := 0;    uiOffset := uiOffset + 1;END_IFIF stPublish.uiPayloadLen > 0 THEN    IF (TO_UDINT(uiOffset) + TO_UDINT(stPublish.uiPayloadLen)) > udiBufferSize THEN        M_BuildPublish := FALSE;        RETURN;    END_IF    // Payload 存在 ST STRING 中时按 0 基下标读取。    // 这和 MQTT 报文字节数组 aBuffer[0..] 的下标体系一致,可以避免转发时首字节丢失、尾部多出垃圾字符。    FOR uiIndex := 0 TO stPublish.uiPayloadLen - 1 DO        aBuffer[uiOffset + uiIndex] := stPublish.sPayload[uiIndex];    END_FOR    uiOffset := uiOffset + stPublish.uiPayloadLen;END_IFuiFrameLen := uiOffset - uiWriteOffset;uiLastFrameLen := uiFrameLen;M_BuildPublish := TRUE;

完整代码 02: FB_MqttBrokerCodec.M_BuildSimpleAck.st

/// =======================================================================/// 名称      : M_BuildSimpleAck/// 功能      : 构建固定长度 MQTT 协议响应/// 说明      : 用于 CONNACK、PUBACK、SUBACK、UNSUBACK、PINGRESP 等轻量响应包。/// 编程人员  : ControlRookie/// 时间      : 2026-05-08/// 版本      : V1.0/// ======================================================================={attribute 'hide_all_locals'}METHOD M_BuildSimpleAck : BOOLVAR_INPUT    ePacketType  : E_MqttPacketType; // 需要构建的 MQTT 响应报文类型    byProtocolLevel : BYTE; // 目标客户端 MQTT 协议级别,5 表示响应中需要携带零属性长度    uiPacketId   : UINT; // 需要带 Packet Identifier 的响应使用该值,无需 PacketId 时为 0    byReturnCode : BYTE; // CONNACK/SUBACK 返回码,其他响应通常为 0    uiReturnCount : UINT; // SUBACK 多 Topic 返回码数量,普通响应传 0    uiWriteOffset : UINT; // 当前响应帧写入发送缓冲区的起始偏移,批量组包时用于追加到上一帧之后[byte]    udiBufferSize : UDINT; // 发送缓冲区总容量[byte]END_VARVAR_IN_OUT    aBuffer      : ARRAY[*] OF BYTE; // MQTT 发送缓冲区    aReturnCodes : ARRAY[*] OF BYTE; // SUBACK 多 Topic 返回码数组,普通响应可传空闲数组END_VARVAR_OUTPUT    uiFrameLen   : UINT; // 构建出的 MQTT 响应报文长度[byte]END_VARVAR    uiIndex      : UINT; // SUBACK 多返回码复制索引[1..cnMaxTopicItemsPerPacket]    udiNeededLen : UDINT; // 当前响应帧需要的完整缓冲长度,包含固定头、可变头和返回码[byte]END_VAR// === IMPLEMENTATION ===uiFrameLen := 0;CASE ePacketType OF    E_MqttPacketType.byConnAck:        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            udiNeededLen := 5;        ELSE            udiNeededLen := 4;        END_IF        IF (TO_UDINT(uiWriteOffset) + udiNeededLen) > udiBufferSize THEN            M_BuildSimpleAck := FALSE;            RETURN;        END_IF        aBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.byConnAck);        aBuffer[uiWriteOffset + 2] := 0;        aBuffer[uiWriteOffset + 3] := byReturnCode;        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            // MQTT 5.0 CONNACK = Acknowledge Flags + Reason Code + Properties。            // 当前轻量兼容层不返回任何属性,因此属性长度固定写 0。            aBuffer[uiWriteOffset + 1] := 3;            aBuffer[uiWriteOffset + 4] := 0;            uiFrameLen := 5;        ELSE            aBuffer[uiWriteOffset + 1] := 2;            uiFrameLen := 4;        END_IF    E_MqttPacketType.byPubAck:        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            udiNeededLen := 6;        ELSE            udiNeededLen := 4;        END_IF        IF (TO_UDINT(uiWriteOffset) + udiNeededLen) > udiBufferSize THEN            M_BuildSimpleAck := FALSE;            RETURN;        END_IF        aBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.byPubAck);        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            aBuffer[uiWriteOffset + 1] := 4;        ELSE            aBuffer[uiWriteOffset + 1] := 2;        END_IF        aBuffer[uiWriteOffset + 2] := TO_BYTE(uiPacketId / 256);        aBuffer[uiWriteOffset + 3] := TO_BYTE(uiPacketId MOD 256);        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            // MQTT 5.0 UNSUBACK = Packet Identifier + Properties + Reason Codes。            // 当前 M_EnqueueProtocolAck 只生成单返回码,属性长度固定写 0。            aBuffer[uiWriteOffset + 4] := 0;            aBuffer[uiWriteOffset + 5] := byReturnCode;            uiFrameLen := 6;        ELSE            uiFrameLen := 4;        END_IF    E_MqttPacketType.byPubRec:        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            udiNeededLen := 6;        ELSE            udiNeededLen := 4;        END_IF        IF (TO_UDINT(uiWriteOffset) + udiNeededLen) > udiBufferSize THEN            M_BuildSimpleAck := FALSE;            RETURN;        END_IF        aBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.byPubRec);        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            aBuffer[uiWriteOffset + 1] := 4;        ELSE            aBuffer[uiWriteOffset + 1] := 2;        END_IF        aBuffer[uiWriteOffset + 2] := TO_BYTE(uiPacketId / 256);        aBuffer[uiWriteOffset + 3] := TO_BYTE(uiPacketId MOD 256);        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            aBuffer[uiWriteOffset + 4] := byReturnCode;            aBuffer[uiWriteOffset + 5] := 0;            uiFrameLen := 6;        ELSE            uiFrameLen := 4;        END_IF    E_MqttPacketType.byPubRel:        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            udiNeededLen := 6;        ELSE            udiNeededLen := 4;        END_IF        IF (TO_UDINT(uiWriteOffset) + udiNeededLen) > udiBufferSize THEN            M_BuildSimpleAck := FALSE;            RETURN;        END_IF        aBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.byPubRel);        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            aBuffer[uiWriteOffset + 1] := 4;        ELSE            aBuffer[uiWriteOffset + 1] := 2;        END_IF        aBuffer[uiWriteOffset + 2] := TO_BYTE(uiPacketId / 256);        aBuffer[uiWriteOffset + 3] := TO_BYTE(uiPacketId MOD 256);        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            aBuffer[uiWriteOffset + 4] := byReturnCode;            aBuffer[uiWriteOffset + 5] := 0;            uiFrameLen := 6;        ELSE            uiFrameLen := 4;        END_IF    E_MqttPacketType.byPubComp:        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            udiNeededLen := 6;        ELSE            udiNeededLen := 4;        END_IF        IF (TO_UDINT(uiWriteOffset) + udiNeededLen) > udiBufferSize THEN            M_BuildSimpleAck := FALSE;            RETURN;        END_IF        aBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.byPubComp);        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            aBuffer[uiWriteOffset + 1] := 4;        ELSE            aBuffer[uiWriteOffset + 1] := 2;        END_IF        aBuffer[uiWriteOffset + 2] := TO_BYTE(uiPacketId / 256);        aBuffer[uiWriteOffset + 3] := TO_BYTE(uiPacketId MOD 256);        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            aBuffer[uiWriteOffset + 4] := byReturnCode;            aBuffer[uiWriteOffset + 5] := 0;            uiFrameLen := 6;        ELSE            uiFrameLen := 4;        END_IF    E_MqttPacketType.bySubAck:        IF uiReturnCount = 0 THEN            IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN                udiNeededLen := 6;            ELSE                udiNeededLen := 5;            END_IF        ELSE            IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN                udiNeededLen := 5 + TO_UDINT(uiReturnCount);            ELSE                udiNeededLen := 4 + TO_UDINT(uiReturnCount);            END_IF        END_IF        IF (TO_UDINT(uiWriteOffset) + udiNeededLen) > udiBufferSize THEN            M_BuildSimpleAck := FALSE;            RETURN;        END_IF        aBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.bySubAck);        IF uiReturnCount = 0 THEN            IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN                aBuffer[uiWriteOffset + 1] := 4;            ELSE                aBuffer[uiWriteOffset + 1] := 3;            END_IF        ELSE            IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN                aBuffer[uiWriteOffset + 1] := TO_BYTE(3 + uiReturnCount);            ELSE                aBuffer[uiWriteOffset + 1] := TO_BYTE(2 + uiReturnCount);            END_IF        END_IF        aBuffer[uiWriteOffset + 2] := TO_BYTE(uiPacketId / 256);        aBuffer[uiWriteOffset + 3] := TO_BYTE(uiPacketId MOD 256);        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            aBuffer[uiWriteOffset + 4] := 0;            IF uiReturnCount = 0 THEN                aBuffer[uiWriteOffset + 5] := byReturnCode;                uiFrameLen := 6;            ELSE                FOR uiIndex := 1 TO uiReturnCount DO                    IF uiIndex > GVL_MqttBroker.cnMaxTopicItemsPerPacket THEN                        M_BuildSimpleAck := FALSE;                        RETURN;                    END_IF                    aBuffer[uiWriteOffset + 4 + uiIndex] := aReturnCodes[uiIndex];                END_FOR                uiFrameLen := 5 + uiReturnCount;            END_IF        ELSE            IF uiReturnCount = 0 THEN                aBuffer[uiWriteOffset + 4] := byReturnCode;                uiFrameLen := 5;            ELSE                FOR uiIndex := 1 TO uiReturnCount DO                    IF uiIndex > GVL_MqttBroker.cnMaxTopicItemsPerPacket THEN                        M_BuildSimpleAck := FALSE;                        RETURN;                    END_IF                    aBuffer[uiWriteOffset + 3 + uiIndex] := aReturnCodes[uiIndex];                END_FOR                uiFrameLen := 4 + uiReturnCount;            END_IF        END_IF    E_MqttPacketType.byUnsubAck:        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            udiNeededLen := 6;        ELSE            udiNeededLen := 4;        END_IF        IF (TO_UDINT(uiWriteOffset) + udiNeededLen) > udiBufferSize THEN            M_BuildSimpleAck := FALSE;            RETURN;        END_IF        aBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.byUnsubAck);        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            aBuffer[uiWriteOffset + 1] := 4;        ELSE            aBuffer[uiWriteOffset + 1] := 2;        END_IF        aBuffer[uiWriteOffset + 2] := TO_BYTE(uiPacketId / 256);        aBuffer[uiWriteOffset + 3] := TO_BYTE(uiPacketId MOD 256);        IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN            aBuffer[uiWriteOffset + 4] := byReturnCode;            aBuffer[uiWriteOffset + 5] := 0;            uiFrameLen := 6;        ELSE            uiFrameLen := 4;        END_IF    E_MqttPacketType.byPingResp:        IF (TO_UDINT(uiWriteOffset) + 2) > udiBufferSize THEN            M_BuildSimpleAck := FALSE;            RETURN;        END_IF        aBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.byPingResp);        aBuffer[uiWriteOffset + 1] := 0;        uiFrameLen := 2;ELSE    M_BuildSimpleAck := FALSE;    RETURN;END_CASEuiLastFrameLen := uiFrameLen;M_BuildSimpleAck := TRUE;

完整代码 03: FB_MqttBrokerCodec.M_ParseConnect.st

/// =======================================================================/// 名称      : M_ParseConnect/// 功能      : 解析 MQTT CONNECT 报文/// 说明      : 校验 MQTT 3.1 / 3.1.1 / 5.0 协议名、协议级别、ClientID,并解析 Clean Session 与 Will。/// 编程人员  : ControlRookie/// 时间      : 2026-05-08/// 版本      : V1.0/// ======================================================================={attribute 'hide_all_locals'}METHOD M_ParseConnect : BOOLVAR_INPUT    uiBodyOffset : UINT; // CONNECT 可变头在报文缓冲区中的起始偏移[byte]    uiFrameLen   : UINT; // 当前 CONNECT 完整报文长度[byte]END_VARVAR_IN_OUT    aBuffer      : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区    stConnection : ST_MqttBrokerConnection; // 当前连接槽位,会被写入 ClientID、KeepAlive 和 Will 信息END_VARVAR_OUTPUT    eError       : E_MqttBrokerError; // 解析失败时返回的 Broker 错误码END_VARVAR    uiOffset       : UINT; // 当前解析偏移[byte]    uiProtocolLen  : UINT; // 协议名字段长度[byte]    uiClientIdLen  : UINT; // ClientID 字段长度[byte]    uiWillTopicLen : UINT; // Will Topic 字段长度[byte]    uiWillMsgLen   : UINT; // Will Payload 字段长度[byte]    uiUsernameLen  : UINT; // 用户名字段长度[byte]    uiPasswordLen  : UINT; // 密码字段长度[byte]    udiPropertyLen : UDINT; // MQTT 5.0 CONNECT 属性区长度,当前轻量兼容模式只校验并跳过[byte]    byFlags        : BYTE; // CONNECT Flags 字节,包含用户名、密码、Will、Clean Session 标志    byProtocolLevel : BYTE; // CONNECT 中声明的 MQTT 协议级别,3 表示 3.1,4 表示 3.1.1,5 表示 5.0    byWillQoS      : BYTE; // 从 CONNECT Flags 中提取的 Will QoS 数值    sTempString    : STRING; // CONNECT 用户名/密码临时缓冲,再按目标字段容量赋值END_VAR// === IMPLEMENTATION ===eError := E_MqttBrokerError.uiNoError;uiOffset := uiBodyOffset;IF uiFrameLen < 14 THEN    eError := E_MqttBrokerError.uiProtocolMalformed;    M_ParseConnect := FALSE;    RETURN;END_IFIF (uiOffset + 1) >= uiFrameLen THEN    eError := E_MqttBrokerError.uiProtocolMalformed;    M_ParseConnect := FALSE;    RETURN;END_IFuiProtocolLen := TO_UINT(aBuffer[uiOffset]) * 256 + TO_UINT(aBuffer[uiOffset + 1]);CASE uiProtocolLen OF    4:        IF (uiOffset + 5) >= uiFrameLen THEN            eError := E_MqttBrokerError.uiProtocolMalformed;            M_ParseConnect := FALSE;            RETURN;        END_IF        IF (aBuffer[uiOffset + 2] <> 16#4D)            OR (aBuffer[uiOffset + 3] <> 16#51)            OR (aBuffer[uiOffset + 4] <> 16#54)            OR (aBuffer[uiOffset + 5] <> 16#54) THEN            eError := E_MqttBrokerError.uiUnsupportedProtocol;            M_ParseConnect := FALSE;            RETURN;        END_IF        uiOffset := uiOffset + 6;    6:        IF (uiOffset + 7) >= uiFrameLen THEN            eError := E_MqttBrokerError.uiProtocolMalformed;            M_ParseConnect := FALSE;            RETURN;        END_IF        // MQTT 3.1 老客户端使用协议名 MQIsdp,协议级别为 3。        // 此处单独分支处理,避免把 3.1.1/5.0 的标准 MQTT 协议名校验放宽。        IF (aBuffer[uiOffset + 2] <> 16#4D)            OR (aBuffer[uiOffset + 3] <> 16#51)            OR (aBuffer[uiOffset + 4] <> 16#49)            OR (aBuffer[uiOffset + 5] <> 16#73)            OR (aBuffer[uiOffset + 6] <> 16#64)            OR (aBuffer[uiOffset + 7] <> 16#70) THEN            eError := E_MqttBrokerError.uiUnsupportedProtocol;            M_ParseConnect := FALSE;            RETURN;        END_IF        uiOffset := uiOffset + 8;ELSE    eError := E_MqttBrokerError.uiUnsupportedProtocol;    M_ParseConnect := FALSE;    RETURN;END_CASEIF uiOffset >= uiFrameLen THEN    eError := E_MqttBrokerError.uiProtocolMalformed;    M_ParseConnect := FALSE;    RETURN;END_IFbyProtocolLevel := aBuffer[uiOffset];IF ((uiProtocolLen = 4)    AND (byProtocolLevel <> GVL_MqttBroker.cnMqttProtocolLevel311)    AND (byProtocolLevel <> GVL_MqttBroker.cnMqttProtocolLevel5))    OR ((uiProtocolLen = 6)    AND (byProtocolLevel <> GVL_MqttBroker.cnMqttProtocolLevel31)) THEN    eError := E_MqttBrokerError.uiUnsupportedProtocol;    M_ParseConnect := FALSE;    RETURN;END_IFstConnection.byProtocolLevel := byProtocolLevel;uiOffset := uiOffset + 1;IF uiOffset >= uiFrameLen THEN    eError := E_MqttBrokerError.uiProtocolMalformed;    M_ParseConnect := FALSE;    RETURN;END_IFbyFlags := aBuffer[uiOffset];uiOffset := uiOffset + 1;IF (byFlags AND 16#01) <> 0 THEN    eError := E_MqttBrokerError.uiProtocolMalformed;    M_ParseConnect := FALSE;    RETURN;END_IFIF (uiOffset + 1) >= uiFrameLen THEN    eError := E_MqttBrokerError.uiProtocolMalformed;    M_ParseConnect := FALSE;    RETURN;END_IFstConnection.uiKeepAlive := TO_UINT(aBuffer[uiOffset]) * 256 + TO_UINT(aBuffer[uiOffset + 1]);IF stConnection.uiKeepAlive = 0 THEN    stConnection.uiKeepAlive := GVL_MqttBroker.cnDefaultKeepAlive;END_IFuiOffset := uiOffset + 2;IF stConnection.byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN    // MQTT 5.0 在 KeepAlive 后增加 CONNECT Properties。    // 当前 Broker 定位为工业轻量兼容,不解释 User Property / Session Expiry 等高级属性,    // 但必须严格跳过属性长度字段,避免后续 ClientID 解析错位导致 5.0 客户端无法连接。    IF NOT F_MqttSkipVariableByteInteger(        aBuffer := aBuffer,        uiOffset := uiOffset,        uiBufferLen := uiFrameLen,        udiValue => udiPropertyLen) THEN        eError := E_MqttBrokerError.uiProtocolMalformed;        M_ParseConnect := FALSE;        RETURN;    END_IF    IF (TO_UDINT(uiOffset) + udiPropertyLen) > TO_UDINT(uiFrameLen) THEN        eError := E_MqttBrokerError.uiProtocolMalformed;        M_ParseConnect := FALSE;        RETURN;    END_IF    uiOffset := uiOffset + TO_UINT(udiPropertyLen);END_IFIF NOT F_MqttReadString(    aBuffer := aBuffer,    uiOffset := uiOffset,    sValue := stConnection.sClientId,    uiBufferLen := uiFrameLen,    uiMaxLen := GVL_MqttBroker.cnMaxClientIdLen,    uiStringLen => uiClientIdLen) THEN    eError := E_MqttBrokerError.uiInvalidClientId;    M_ParseConnect := FALSE;    RETURN;END_IFIF uiClientIdLen = 0 THEN    eError := E_MqttBrokerError.uiInvalidClientId;    M_ParseConnect := FALSE;    RETURN;END_IFstConnection.xCleanSession := (byFlags AND 16#02) <> 0;stConnection.xWillFlag := (byFlags AND 16#04) <> 0;stConnection.xWillRetain := (byFlags AND 16#20) <> 0;stConnection.sUsername := '';stConnection.sPassword := '';stConnection.xAuthenticated := FALSE;byWillQoS := SHR(byFlags AND 16#183);CASE byWillQoS OF    0:        stConnection.eWillQoS := E_MqttQoS.byQoS0;    1:        stConnection.eWillQoS := E_MqttQoS.byQoS1;    2:        stConnection.eWillQoS := E_MqttQoS.byQoS2;ELSE    eError := E_MqttBrokerError.uiProtocolMalformed;    M_ParseConnect := FALSE;    RETURN;END_CASEIF stConnection.xWillFlag THEN    IF stConnection.byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN        // MQTT 5.0 Will Properties 位于 Will Topic 之前。        // 当前只支持 Will Topic/Payload/QoS/Retain 主链路,属性区做长度校验后跳过。        IF NOT F_MqttSkipVariableByteInteger(            aBuffer := aBuffer,            uiOffset := uiOffset,            uiBufferLen := uiFrameLen,            udiValue => udiPropertyLen) THEN            eError := E_MqttBrokerError.uiProtocolMalformed;            M_ParseConnect := FALSE;            RETURN;        END_IF        IF (TO_UDINT(uiOffset) + udiPropertyLen) > TO_UDINT(uiFrameLen) THEN            eError := E_MqttBrokerError.uiProtocolMalformed;            M_ParseConnect := FALSE;            RETURN;        END_IF        uiOffset := uiOffset + TO_UINT(udiPropertyLen);    END_IF    IF NOT F_MqttReadString(        aBuffer := aBuffer,        uiOffset := uiOffset,        sValue := stConnection.sWillTopic,        uiBufferLen := uiFrameLen,        uiMaxLen := GVL_MqttBroker.cnMaxTopicLen,        uiStringLen => uiWillTopicLen) THEN        eError := E_MqttBrokerError.uiInvalidTopic;        M_ParseConnect := FALSE;        RETURN;    END_IF    IF NOT F_MqttIsValidTopicName(sTopic := stConnection.sWillTopic) THEN        eError := E_MqttBrokerError.uiInvalidTopic;        M_ParseConnect := FALSE;        RETURN;    END_IF    IF NOT F_MqttReadString(        aBuffer := aBuffer,        uiOffset := uiOffset,        sValue := stConnection.sWillPayload,        uiBufferLen := uiFrameLen,        uiMaxLen := GVL_MqttBroker.cnMaxPayloadLen,        uiStringLen => uiWillMsgLen) THEN        eError := E_MqttBrokerError.uiProtocolMalformed;        M_ParseConnect := FALSE;        RETURN;    END_IFEND_IFIF (byFlags AND 16#80) <> 0 THEN    IF NOT F_MqttReadString(        aBuffer := aBuffer,        uiOffset := uiOffset,        sValue := sTempString,        uiBufferLen := uiFrameLen,        uiMaxLen := GVL_MqttBroker.cnMaxUsernameLen,        uiStringLen => uiUsernameLen) THEN        eError := E_MqttBrokerError.uiProtocolMalformed;        M_ParseConnect := FALSE;        RETURN;    END_IF    stConnection.sUsername := sTempString;END_IFIF (byFlags AND 16#40) <> 0 THEN    IF NOT F_MqttReadString(        aBuffer := aBuffer,        uiOffset := uiOffset,        sValue := sTempString,        uiBufferLen := uiFrameLen,        uiMaxLen := GVL_MqttBroker.cnMaxPasswordLen,        uiStringLen => uiPasswordLen) THEN        eError := E_MqttBrokerError.uiProtocolMalformed;        M_ParseConnect := FALSE;        RETURN;    END_IF    stConnection.sPassword := sTempString;END_IFuiLastFrameLen := uiFrameLen;M_ParseConnect := TRUE;

完整代码 04: FB_MqttBrokerCodec.M_ParsePublish.st

/// =======================================================================/// 名称      : M_ParsePublish/// 功能      : 解析 MQTT PUBLISH 报文/// 说明      : 提取 Topic、Payload、QoS、Retain、DUP 和 PacketId,供路由层使用。/// 编程人员  : ControlRookie/// 时间      : 2026-05-08/// 版本      : V1.0/// ======================================================================={attribute 'hide_all_locals'}METHOD M_ParsePublish : BOOLVAR_INPUT    uiSourceSlot : UINT; // 发布来源客户端槽位编号[1..cnMaxClientSlots]    byProtocolLevel : BYTE; // 当前连接协商的 MQTT 协议级别,5 表示需要跳过 PUBLISH Properties    uiBodyOffset : UINT; // PUBLISH 可变头起始偏移[byte]    uiFrameLen   : UINT; // 当前 PUBLISH 完整报文长度[byte]END_VARVAR_IN_OUT    aBuffer      : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区    stPublish    : ST_MqttBrokerPublishFrame; // 解析后的标准发布帧END_VARVAR_OUTPUT    eError       : E_MqttBrokerError; // 解析失败时返回的 Broker 错误码END_VARVAR    uiOffset     : UINT; // 当前解析偏移[byte]    uiTopicLen   : UINT; // Topic Name 字段长度[byte]    uiPayloadIdx : UINT; // Payload 字节复制索引[byte]    udiPropertyLen : UDINT; // MQTT 5.0 PUBLISH 属性区长度,当前轻量兼容模式只校验并跳过[byte]    byQoS        : BYTE; // 固定头中解析出的 QoS 数值END_VAR// === IMPLEMENTATION ===eError := E_MqttBrokerError.uiNoError;stPublish.xValid := FALSE;stPublish.uiSourceSlot := uiSourceSlot;stPublish.uiTargetSlot := 0;stPublish.uiPacketId := 0;stPublish.sTopic := '';stPublish.sPayload := '';uiOffset := uiBodyOffset;IF uiFrameLen <= uiBodyOffset THEN    eError := E_MqttBrokerError.uiProtocolMalformed;    M_ParsePublish := FALSE;    RETURN;END_IFstPublish.xDup := (aBuffer[0] AND 16#08) <> 0;stPublish.xRetain := (aBuffer[0] AND 16#01) <> 0;byQoS := SHR(aBuffer[0] AND 16#061);CASE byQoS OF    0:        stPublish.eQoS := E_MqttQoS.byQoS0;    1:        stPublish.eQoS := E_MqttQoS.byQoS1;    2:        stPublish.eQoS := E_MqttQoS.byQoS2;ELSE    eError := E_MqttBrokerError.uiProtocolMalformed;    M_ParsePublish := FALSE;    RETURN;END_CASEIF NOT F_MqttReadString(    aBuffer := aBuffer,    uiOffset := uiOffset,    sValue := stPublish.sTopic,    uiBufferLen := uiFrameLen,    uiMaxLen := GVL_MqttBroker.cnMaxTopicLen,    uiStringLen => uiTopicLen) THEN    eError := E_MqttBrokerError.uiInvalidTopic;    M_ParsePublish := FALSE;    RETURN;END_IFIF NOT F_MqttIsValidTopicName(sTopic := stPublish.sTopic) THEN    eError := E_MqttBrokerError.uiInvalidTopic;    M_ParsePublish := FALSE;    RETURN;END_IFstPublish.uiTopicLen := uiTopicLen;IF stPublish.eQoS <> E_MqttQoS.byQoS0 THEN    IF (uiOffset + 1) >= uiFrameLen THEN        eError := E_MqttBrokerError.uiProtocolMalformed;        M_ParsePublish := FALSE;        RETURN;    END_IF    stPublish.uiPacketId := TO_UINT(aBuffer[uiOffset]) * 256 + TO_UINT(aBuffer[uiOffset + 1]);    uiOffset := uiOffset + 2;    IF stPublish.uiPacketId = 0 THEN        eError := E_MqttBrokerError.uiProtocolMalformed;        M_ParsePublish := FALSE;        RETURN;    END_IFEND_IFIF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN    // MQTT 5.0 在 PUBLISH 可变头末尾增加 Properties。    // 当前 Broker 不实现 Topic Alias、Payload Format 等 5.0 高级属性,    // 但必须跳过属性区,否则 Payload 会被错误地带上属性长度字节。    IF NOT F_MqttSkipVariableByteInteger(        aBuffer := aBuffer,        uiOffset := uiOffset,        uiBufferLen := uiFrameLen,        udiValue => udiPropertyLen) THEN        eError := E_MqttBrokerError.uiProtocolMalformed;        M_ParsePublish := FALSE;        RETURN;    END_IF    IF (TO_UDINT(uiOffset) + udiPropertyLen) > TO_UDINT(uiFrameLen) THEN        eError := E_MqttBrokerError.uiProtocolMalformed;        M_ParsePublish := FALSE;        RETURN;    END_IF    uiOffset := uiOffset + TO_UINT(udiPropertyLen);END_IFstPublish.uiPayloadLen := uiFrameLen - uiOffset;IF stPublish.uiPayloadLen > GVL_MqttBroker.cnMaxPayloadLen THEN    eError := E_MqttBrokerError.uiPacketTooLarge;    M_ParsePublish := FALSE;    RETURN;END_IFIF stPublish.uiPayloadLen > 0 THEN    // CodeSys/CODESYS 的 STRING 字符访问按 0 基下标工作。    // 这里保存 Payload 时必须与 F_MqttReadString/F_MqttAppendString 保持一致;    // 否则应用层字符串会整体错位,后续转发给订阅客户端时可能构造出非法 PUBLISH。    FOR uiPayloadIdx := 0 TO stPublish.uiPayloadLen - 1 DO        stPublish.sPayload[uiPayloadIdx] := aBuffer[uiOffset + uiPayloadIdx];    END_FOR    stPublish.sPayload[stPublish.uiPayloadLen] := 0;END_IFstPublish.xValid := TRUE;uiLastFrameLen := uiFrameLen;M_ParsePublish := TRUE;

完整代码 05: FB_MqttBrokerCodec.M_ParseSubscribe.st

/// =======================================================================/// 名称      : M_ParseSubscribe/// 功能      : 解析 MQTT SUBSCRIBE 报文/// 说明      : 第二阶段支持单个 SUBSCRIBE 报文内多个 Topic Filter,并逐项输出返回码。/// 编程人员  : ControlRookie/// 时间      : 2026-05-08/// 版本      : V1.0/// ======================================================================={attribute 'hide_all_locals'}METHOD M_ParseSubscribe : BOOLVAR_INPUT    byProtocolLevel : BYTE; // 当前连接协商的 MQTT 协议级别,5 表示需要跳过 SUBSCRIBE Properties    uiBodyOffset : UINT; // SUBSCRIBE 可变头起始偏移[byte]    uiFrameLen   : UINT; // 当前 SUBSCRIBE 完整报文长度[byte]END_VARVAR_IN_OUT    aBuffer      : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区    aTopicItems  : ARRAY[*] OF ST_MqttBrokerTopicItem; // 解析出的多 Topic 订阅条目数组END_VARVAR_OUTPUT    uiPacketId   : UINT; // SUBSCRIBE Packet Identifier    uiItemCount  : UINT; // 本次 SUBSCRIBE 成功解析出的 Topic 条目数量    eError       : E_MqttBrokerError; // 解析失败时返回的 Broker 错误码END_VARVAR    uiOffset     : UINT; // 当前解析偏移[byte]    uiFilterLen  : UINT; // Topic Filter 字段长度[byte]    uiIndex      : UINT; // Topic 条目数组写入索引[1..cnMaxTopicItemsPerPacket]    udiPropertyLen : UDINT; // MQTT 5.0 SUBSCRIBE 属性区长度,当前轻量兼容模式只校验并跳过[byte]    byQoS        : BYTE; // SUBSCRIBE 载荷中请求的 QoS 数值END_VAR// === IMPLEMENTATION ===eError := E_MqttBrokerError.uiNoError;uiPacketId := 0;uiItemCount := 0;uiOffset := uiBodyOffset;FOR uiIndex := 1 TO GVL_MqttBroker.cnMaxTopicItemsPerPacket DO    aTopicItems[uiIndex].xUsed := FALSE;    aTopicItems[uiIndex].sTopicFilter := '';    aTopicItems[uiIndex].uiFilterLen := 0;    aTopicItems[uiIndex].eQoS := E_MqttQoS.byQoS0;    aTopicItems[uiIndex].byReturnCode := 16#80;END_FORIF (uiOffset + 1) >= uiFrameLen THEN    eError := E_MqttBrokerError.uiProtocolMalformed;    M_ParseSubscribe := FALSE;    RETURN;END_IFuiPacketId := TO_UINT(aBuffer[uiOffset]) * 256 + TO_UINT(aBuffer[uiOffset + 1]);uiOffset := uiOffset + 2;IF uiPacketId = 0 THEN    eError := E_MqttBrokerError.uiProtocolMalformed;    M_ParseSubscribe := FALSE;    RETURN;END_IFIF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN    // MQTT 5.0 在 SUBSCRIBE Packet Identifier 后增加 Properties。    // 当前轻量 Broker 不解释 Subscription Identifier / User Property 等高级属性,    // 只跳过属性区,让后续 Topic Filter 列表按正确偏移解析。    IF NOT F_MqttSkipVariableByteInteger(        aBuffer := aBuffer,        uiOffset := uiOffset,        uiBufferLen := uiFrameLen,        udiValue => udiPropertyLen) THEN        eError := E_MqttBrokerError.uiProtocolMalformed;        M_ParseSubscribe := FALSE;        RETURN;    END_IF    IF (TO_UDINT(uiOffset) + udiPropertyLen) > TO_UDINT(uiFrameLen) THEN        eError := E_MqttBrokerError.uiProtocolMalformed;        M_ParseSubscribe := FALSE;        RETURN;    END_IF    uiOffset := uiOffset + TO_UINT(udiPropertyLen);END_IFWHILE uiOffset < uiFrameLen DO    IF uiItemCount >= GVL_MqttBroker.cnMaxTopicItemsPerPacket THEN        eError := E_MqttBrokerError.uiPacketTooLarge;        M_ParseSubscribe := FALSE;        RETURN;    END_IF    uiIndex := uiItemCount + 1;    IF NOT F_MqttReadString(        aBuffer := aBuffer,        uiOffset := uiOffset,        sValue := aTopicItems[uiIndex].sTopicFilter,        uiBufferLen := uiFrameLen,        uiMaxLen := GVL_MqttBroker.cnMaxTopicLen,        uiStringLen => uiFilterLen) THEN        eError := E_MqttBrokerError.uiInvalidTopic;        M_ParseSubscribe := FALSE;        RETURN;    END_IF    IF uiOffset >= uiFrameLen THEN        eError := E_MqttBrokerError.uiProtocolMalformed;        M_ParseSubscribe := FALSE;        RETURN;    END_IF    byQoS := aBuffer[uiOffset];    uiOffset := uiOffset + 1;    CASE byQoS OF        0:            aTopicItems[uiIndex].eQoS := E_MqttQoS.byQoS0;        1:            aTopicItems[uiIndex].eQoS := E_MqttQoS.byQoS1;        2:            aTopicItems[uiIndex].eQoS := E_MqttQoS.byQoS2;    ELSE        eError := E_MqttBrokerError.uiUnsupportedQoS;        M_ParseSubscribe := FALSE;        RETURN;    END_CASE    IF NOT F_MqttIsValidTopicFilter(sFilter := aTopicItems[uiIndex].sTopicFilter) THEN        eError := E_MqttBrokerError.uiInvalidTopic;        M_ParseSubscribe := FALSE;        RETURN;    END_IF    aTopicItems[uiIndex].xUsed := TRUE;    aTopicItems[uiIndex].uiFilterLen := uiFilterLen;    aTopicItems[uiIndex].byReturnCode := TO_BYTE(aTopicItems[uiIndex].eQoS);    uiItemCount := uiItemCount + 1;END_WHILEIF uiItemCount = 0 THEN    eError := E_MqttBrokerError.uiProtocolMalformed;    M_ParseSubscribe := FALSE;    RETURN;END_IFM_ParseSubscribe := TRUE;

完整代码 06: FB_MqttBrokerCodec.M_ParseUnsubscribe.st

/// =======================================================================/// 名称      : M_ParseUnsubscribe/// 功能      : 解析 MQTT UNSUBSCRIBE 报文/// 说明      : 第二阶段支持单个 UNSUBSCRIBE 报文内多个 Topic Filter。/// 编程人员  : ControlRookie/// 时间      : 2026-05-08/// 版本      : V1.0/// ======================================================================={attribute 'hide_all_locals'}METHOD M_ParseUnsubscribe : BOOLVAR_INPUT    byProtocolLevel : BYTE; // 当前连接协商的 MQTT 协议级别,5 表示需要跳过 UNSUBSCRIBE Properties    uiBodyOffset : UINT; // UNSUBSCRIBE 可变头起始偏移[byte]    uiFrameLen   : UINT; // 当前 UNSUBSCRIBE 完整报文长度[byte]END_VARVAR_IN_OUT    aBuffer      : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区    aTopicItems  : ARRAY[*] OF ST_MqttBrokerTopicItem; // 解析出的多 Topic 取消订阅条目数组END_VARVAR_OUTPUT    uiPacketId   : UINT; // UNSUBSCRIBE Packet Identifier    uiItemCount  : UINT; // 本次 UNSUBSCRIBE 成功解析出的 Topic 条目数量    eError       : E_MqttBrokerError; // 解析失败时返回的 Broker 错误码END_VARVAR    uiOffset     : UINT; // 当前解析偏移[byte]    uiFilterLen  : UINT; // Topic Filter 字段长度[byte]    uiIndex      : UINT; // Topic 条目数组写入索引[1..cnMaxTopicItemsPerPacket]    udiPropertyLen : UDINT; // MQTT 5.0 UNSUBSCRIBE 属性区长度,当前轻量兼容模式只校验并跳过[byte]END_VAR// === IMPLEMENTATION ===eError := E_MqttBrokerError.uiNoError;uiPacketId := 0;uiItemCount := 0;uiOffset := uiBodyOffset;FOR uiIndex := 1 TO GVL_MqttBroker.cnMaxTopicItemsPerPacket DO    aTopicItems[uiIndex].xUsed := FALSE;    aTopicItems[uiIndex].sTopicFilter := '';    aTopicItems[uiIndex].uiFilterLen := 0;    aTopicItems[uiIndex].eQoS := E_MqttQoS.byQoS0;    aTopicItems[uiIndex].byReturnCode := 0;END_FORIF (uiOffset + 1) >= uiFrameLen THEN    eError := E_MqttBrokerError.uiProtocolMalformed;    M_ParseUnsubscribe := FALSE;    RETURN;END_IFuiPacketId := TO_UINT(aBuffer[uiOffset]) * 256 + TO_UINT(aBuffer[uiOffset + 1]);uiOffset := uiOffset + 2;IF uiPacketId = 0 THEN    eError := E_MqttBrokerError.uiProtocolMalformed;    M_ParseUnsubscribe := FALSE;    RETURN;END_IFIF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN    // MQTT 5.0 在 UNSUBSCRIBE Packet Identifier 后增加 Properties。    // 当前只支持 Topic Filter 主链路,属性区做长度校验后跳过。    IF NOT F_MqttSkipVariableByteInteger(        aBuffer := aBuffer,        uiOffset := uiOffset,        uiBufferLen := uiFrameLen,        udiValue => udiPropertyLen) THEN        eError := E_MqttBrokerError.uiProtocolMalformed;        M_ParseUnsubscribe := FALSE;        RETURN;    END_IF    IF (TO_UDINT(uiOffset) + udiPropertyLen) > TO_UDINT(uiFrameLen) THEN        eError := E_MqttBrokerError.uiProtocolMalformed;        M_ParseUnsubscribe := FALSE;        RETURN;    END_IF    uiOffset := uiOffset + TO_UINT(udiPropertyLen);END_IFWHILE uiOffset < uiFrameLen DO    IF uiItemCount >= GVL_MqttBroker.cnMaxTopicItemsPerPacket THEN        eError := E_MqttBrokerError.uiPacketTooLarge;        M_ParseUnsubscribe := FALSE;        RETURN;    END_IF    uiIndex := uiItemCount + 1;    IF NOT F_MqttReadString(        aBuffer := aBuffer,        uiOffset := uiOffset,        sValue := aTopicItems[uiIndex].sTopicFilter,        uiBufferLen := uiFrameLen,        uiMaxLen := GVL_MqttBroker.cnMaxTopicLen,        uiStringLen => uiFilterLen) THEN        eError := E_MqttBrokerError.uiInvalidTopic;        M_ParseUnsubscribe := FALSE;        RETURN;    END_IF    IF NOT F_MqttIsValidTopicFilter(sFilter := aTopicItems[uiIndex].sTopicFilter) THEN        eError := E_MqttBrokerError.uiInvalidTopic;        M_ParseUnsubscribe := FALSE;        RETURN;    END_IF    aTopicItems[uiIndex].xUsed := TRUE;    aTopicItems[uiIndex].uiFilterLen := uiFilterLen;    aTopicItems[uiIndex].eQoS := E_MqttQoS.byQoS0;    aTopicItems[uiIndex].byReturnCode := 0;    uiItemCount := uiItemCount + 1;END_WHILEIF uiItemCount = 0 THEN    eError := E_MqttBrokerError.uiProtocolMalformed;    M_ParseUnsubscribe := FALSE;    RETURN;END_IFM_ParseUnsubscribe := TRUE;

完整代码 07: FB_MqttBrokerCodec.st

/// =======================================================================/// 名称      : FB_MqttBrokerCodec/// 功能      : MQTT Broker 报文编解码器/// 说明      : 本功能块不保存业务状态,只通过方法解析入站报文和构建服务端回包。/// 编程人员  : ControlRookie/// 时间      : 2026-05-08/// 版本      : V1.0/// ======================================================================={attribute 'hide_all_locals'}FUNCTION_BLOCK FB_MqttBrokerCodecVAR    uiLastFrameLen : UINT// 最近一次成功构建或解析的 MQTT 报文长度[byte]END_VAR// === IMPLEMENTATION ===

完整代码 08: F_MqttAppendString.st

/// =======================================================================/// 名称      : F_MqttAppendString/// 功能      : 向 MQTT 报文缓冲区追加 UTF-8 字符串字段/// 说明      : MQTT 字符串格式为 2 字节大端长度 + 字符串内容,本函数负责边界保护。/// 编程人员  : ControlRookie/// 时间      : 2026-05-08/// 版本      : V1.0/// ======================================================================={attribute 'hide_all_locals'}FUNCTION F_MqttAppendString : BOOLVAR_INPUT    sValue        : STRING; // 需要追加到 MQTT 报文中的字符串内容    udiBufferSize : UDINT; // 调用方报文缓冲区总容量[byte]END_VARVAR_IN_OUT    aBuffer       : ARRAY[*] OF BYTE; // 调用方报文缓冲区    uiOffset      : UINT; // 当前写入偏移,成功后推进到字符串末尾后一字节[byte]END_VARVAR    uiLen         : UINT; // 字符串长度[byte]    uiIndex       : UINT; // 字符串字符复制索引,CODESYS STRING 字符下标从 0 开始[byte]    uiWriteIndex  : UINT; // 当前写入报文缓冲区索引[byte]END_VAR// === IMPLEMENTATION ===uiLen := TO_UINT(LEN(sValue));IF (TO_UDINT(uiOffset) + 2 + TO_UDINT(uiLen)) > udiBufferSize THEN    F_MqttAppendString := FALSE;    RETURN;END_IFaBuffer[uiOffset] := TO_BYTE(uiLen / 256);aBuffer[uiOffset + 1] := TO_BYTE(uiLen MOD 256);uiOffset := uiOffset + 2;IF uiLen > 0 THEN    FOR uiIndex := 0 TO uiLen - 1 DO        uiWriteIndex := uiOffset + uiIndex;        IF TO_UDINT(uiWriteIndex) >= udiBufferSize THEN            F_MqttAppendString := FALSE;            RETURN;        END_IF        aBuffer[uiWriteIndex] := TO_BYTE(sValue[uiIndex]);    END_FOREND_IFuiOffset := uiOffset + uiLen;F_MqttAppendString := TRUE;

完整代码 09: F_MqttContainsWildcard.st

/// =======================================================================/// 名称      : F_MqttContainsWildcard/// 功能      : 判断 Topic Filter 是否包含 MQTT 通配符/// 说明      : 供 ACL 判断是否允许 + 或 # 通配符订阅。/// 编程人员  : ControlRookie/// 时间      : 2026-05-08/// 版本      : V1.0/// ======================================================================={attribute 'hide_all_locals'}FUNCTION F_MqttContainsWildcard : BOOLVAR_INPUT    sTopicFilter : STRING; // 需要检查的 Topic FilterEND_VARVAR    uiLen   : UINT// Topic Filter 长度[byte]    uiIndex : UINT// 字符逐字节扫描索引[byte]END_VAR// === IMPLEMENTATION ===uiLen := TO_UINT(LEN(sTopicFilter));IF uiLen = 0 THEN    F_MqttContainsWildcard := FALSE;    RETURN;END_IFFOR uiIndex := 0 TO uiLen - 1 DO    IF (sTopicFilter[uiIndex] = 16#2B) OR (sTopicFilter[uiIndex] = 16#23) THEN        F_MqttContainsWildcard := TRUE;        RETURN;    END_IFEND_FORF_MqttContainsWildcard := FALSE;

完整代码 10: F_MqttDecodeRemainingLength.st

/// =======================================================================/// 名称      : F_MqttDecodeRemainingLength/// 功能      : 解码 MQTT Remaining Length/// 说明      : 从固定报头第 2 字节开始解析 MQTT 变长整数,最多读取 4 字节。/// 编程人员  : ControlRookie/// 时间      : 2026-05-08/// 版本      : V1.0/// ======================================================================={attribute 'hide_all_locals'}FUNCTION F_MqttDecodeRemainingLength : BOOLVAR_INPUT    uiStartIndex : UINT; // Remaining Length 在缓冲区中的起始索引,通常为 1[byte]    uiBufferLen  : UINT; // 当前缓冲区内有效数据长度[byte]END_VARVAR_IN_OUT    aBuffer      : ARRAY[*] OF BYTE; // MQTT 原始接收缓冲区END_VARVAR_OUTPUT    udiValue     : UDINT; // 解码后的 Remaining Length 数值[byte]    uiBytesUsed  : UINT; // Remaining Length 字段实际占用的字节数[byte]    xNeedMore    : BOOL; // 数据不足时置 TRUE,上层应继续读取 TCP 数据END_VARVAR    udiMultiplier : UDINT; // MQTT 变长整数倍率,依次为 1、128、16384、2097152    byEncoded     : BYTE; // 当前读取的编码字节    uiIndex       : UINT; // 当前读取缓冲区索引[byte]    uiLoop        : UINT; // 变长整数最多 4 字节的循环计数END_VAR// === IMPLEMENTATION ===udiValue := 0;uiBytesUsed := 0;xNeedMore := FALSE;udiMultiplier := 1;uiIndex := uiStartIndex;IF uiStartIndex >= uiBufferLen THEN    xNeedMore := TRUE;    F_MqttDecodeRemainingLength := FALSE;    RETURN;END_IFFOR uiLoop := 1 TO 4 DO    IF uiIndex >= uiBufferLen THEN        xNeedMore := TRUE;        F_MqttDecodeRemainingLength := FALSE;        RETURN;    END_IF    byEncoded := aBuffer[uiIndex];    udiValue := udiValue + TO_UDINT(byEncoded AND 16#7F) * udiMultiplier;    uiBytesUsed := uiBytesUsed + 1;    IF (byEncoded AND 16#80) = 0 THEN        F_MqttDecodeRemainingLength := TRUE;        RETURN;    END_IF    udiMultiplier := udiMultiplier * 128;    uiIndex := uiIndex + 1;END_FORF_MqttDecodeRemainingLength := FALSE;

完整代码 11: F_MqttEncodeRemainingLength.st

/// =======================================================================/// 名称      : F_MqttEncodeRemainingLength/// 功能      : 编码 MQTT Remaining Length/// 说明      : 把 MQTT 剩余长度编码为 1~4 字节变长整数,并写入调用方提供的缓冲区。/// 编程人员  : ControlRookie/// 时间      : 2026-05-08/// 版本      : V1.0/// ======================================================================={attribute 'hide_all_locals'}FUNCTION F_MqttEncodeRemainingLength : BOOLVAR_INPUT    udiValue     : UDINT; // 需要编码的 MQTT Remaining Length 数值[byte]    udiBufferSize : UDINT; // 调用方输出缓冲区容量[byte]END_VARVAR_IN_OUT    aBuffer      : ARRAY[*] OF BYTE; // 调用方输出缓冲区,函数从索引 0 起写入编码结果END_VARVAR_OUTPUT    uiEncodedLen : UINT; // 实际编码产生的字节数[byte]END_VARVAR    udiWorkValue : UDINT; // 编码过程中逐步除以 128 的临时值    byEncoded    : BYTE; // 当前轮生成的 7 位数据和 continuation 标志    uiIndex      : UINT; // 当前写入缓冲区的索引[byte]END_VAR// === IMPLEMENTATION ===uiEncodedLen := 0;udiWorkValue := udiValue;IF udiValue > GVL_MqttBroker.cnMqttRemainingLengthMax THEN    F_MqttEncodeRemainingLength := FALSE;    RETURN;END_IFREPEAT    IF TO_UDINT(uiIndex) >= udiBufferSize THEN        F_MqttEncodeRemainingLength := FALSE;        RETURN;    END_IF    byEncoded := TO_BYTE(udiWorkValue MOD 128);    udiWorkValue := udiWorkValue / 128;    IF udiWorkValue > 0 THEN        byEncoded := byEncoded OR 16#80;    END_IF    aBuffer[uiIndex] := byEncoded;    uiIndex := uiIndex + 1;UNTIL udiWorkValue = 0END_REPEATuiEncodedLen := uiIndex;F_MqttEncodeRemainingLength := TRUE;

完整代码 12: F_MqttIsValidTopicFilter.st

/// =======================================================================/// 名称      : F_MqttIsValidTopicFilter/// 功能      : 校验 MQTT Topic Filter/// 说明      : SUBSCRIBE 使用的过滤器允许 + / #,但必须满足 MQTT 通配符位置规则。/// 编程人员  : ControlRookie/// 时间      : 2026-05-08/// 版本      : V1.0/// ======================================================================={attribute 'hide_all_locals'}FUNCTION F_MqttIsValidTopicFilter : BOOLVAR_INPUT    sFilter : STRING; // 待校验的 MQTT Topic FilterEND_VARVAR    uiLen      : UINT; // Topic Filter 长度[byte]    uiIndex    : UINT; // 当前检查的字符位置,CODESYS STRING 字符下标从 0 开始[byte]    byPrev     : BYTE; // 当前字符前一个字符,用于判断通配符是否独占层级    byNext     : BYTE; // 当前字符后一个字符,用于判断通配符是否独占层级END_VAR// === IMPLEMENTATION ===uiLen := TO_UINT(LEN(sFilter));IF uiLen = 0 THEN    F_MqttIsValidTopicFilter := FALSE;    RETURN;END_IFIF uiLen > GVL_MqttBroker.cnMaxTopicLen THEN    F_MqttIsValidTopicFilter := FALSE;    RETURN;END_IFFOR uiIndex := 0 TO uiLen - 1 DO    IF uiIndex > 0 THEN        byPrev := sFilter[uiIndex - 1];    ELSE        byPrev := 0;    END_IF    IF uiIndex < (uiLen - 1) THEN        byNext := sFilter[uiIndex + 1];    ELSE        byNext := 0;    END_IF    CASE sFilter[uiIndex] OF        16#23:            // # 必须是最后一个字符,并且要么单独出现,要么前面是层级分隔符 /。            IF uiIndex <> (uiLen - 1) THEN                F_MqttIsValidTopicFilter := FALSE;                RETURN;            END_IF            IF (uiIndex > 0) AND (byPrev <> 16#2F) THEN                F_MqttIsValidTopicFilter := FALSE;                RETURN;            END_IF        16#2B:            // + 必须独占一个层级,左右只能是边界或层级分隔符 /。            IF (uiIndex > 0) AND (byPrev <> 16#2F) THEN                F_MqttIsValidTopicFilter := FALSE;                RETURN;            END_IF            IF (uiIndex < (uiLen - 1)) AND (byNext <> 16#2F) THEN                F_MqttIsValidTopicFilter := FALSE;                RETURN;            END_IF    ELSE        // 普通字符不需要额外限制;UTF-8 合法性由上位系统或客户端侧保证。    END_CASEEND_FORF_MqttIsValidTopicFilter := TRUE;

完整代码 13: F_MqttIsValidTopicName.st

/// =======================================================================/// 名称      : F_MqttIsValidTopicName/// 功能      : 校验 MQTT Topic Name/// 说明      : PUBLISH 使用的 Topic Name 不能为空,且不能包含 + / # 通配符。/// 编程人员  : ControlRookie/// 时间      : 2026-05-08/// 版本      : V1.0/// ======================================================================={attribute 'hide_all_locals'}FUNCTION F_MqttIsValidTopicName : BOOLVAR_INPUT    sTopic : STRING; // 待校验的 MQTT Topic NameEND_VARVAR    uiLen   : UINT// Topic Name 长度[byte]    uiIndex : UINT// 当前检查的字符位置,CODESYS STRING 字符下标从 0 开始[byte]END_VAR// === IMPLEMENTATION ===uiLen := TO_UINT(LEN(sTopic));IF uiLen = 0 THEN    F_MqttIsValidTopicName := FALSE;    RETURN;END_IFIF uiLen > GVL_MqttBroker.cnMaxTopicLen THEN    F_MqttIsValidTopicName := FALSE;    RETURN;END_IFFOR uiIndex := 0 TO uiLen - 1 DO    IF (sTopic[uiIndex] = 16#2B) OR (sTopic[uiIndex] = 16#23) THEN        F_MqttIsValidTopicName := FALSE;        RETURN;    END_IFEND_FORF_MqttIsValidTopicName := TRUE;

完整代码 14: F_MqttReadString.st

/// =======================================================================/// 名称      : F_MqttReadString/// 功能      : 从 MQTT 报文缓冲区读取 UTF-8 字符串字段/// 说明      : 读取 2 字节大端长度和后续内容,并推进调用方偏移。/// 注意      : CodeSys/CODESYS STRING 字符下标按 0 基访问,不能按 1 基复制。/// 编程人员  : ControlRookie/// 时间      : 2026-05-08/// 版本      : V1.0/// ======================================================================={attribute 'hide_all_locals'}FUNCTION F_MqttReadString : BOOLVAR_INPUT    uiBufferLen : UINT; // 当前 MQTT 报文有效长度[byte]    uiMaxLen    : UINT; // 输出字符串允许保存的最大长度[byte]END_VARVAR_IN_OUT    aBuffer     : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区    uiOffset    : UINT; // 当前读取偏移,成功后推进到字符串末尾后一字节[byte]    sValue      : STRING; // 读取出的字符串内容,短字段由调用方使用默认长度临时缓冲转接END_VARVAR_OUTPUT    uiStringLen : UINT; // MQTT 字符串字段声明的原始长度[byte]END_VARVAR    uiIndex     : UINT; // 字符串字符复制索引,CODESYS STRING 字符下标从 0 开始[byte]    uiReadIndex : UINT; // 当前读取报文缓冲区索引[byte]END_VAR// === IMPLEMENTATION ===sValue := '';uiStringLen := 0;IF (uiOffset + 1) >= uiBufferLen THEN    F_MqttReadString := FALSE;    RETURN;END_IFuiStringLen := TO_UINT(aBuffer[uiOffset]) * 256 + TO_UINT(aBuffer[uiOffset + 1]);uiOffset := uiOffset + 2;IF uiStringLen > uiMaxLen THEN    F_MqttReadString := FALSE;    RETURN;END_IFIF (TO_UDINT(uiOffset) + TO_UDINT(uiStringLen)) > TO_UDINT(uiBufferLen) THEN    F_MqttReadString := FALSE;    RETURN;END_IFIF uiStringLen > 0 THEN    // 关键坑位:    // MQTT 报文缓冲区本身是 ARRAY[0..],CodeSys/CODESYS 的 STRING 字符访问同样按 0 基下标工作。    // 曾经按 1 基写入会导致 CONNECT 协议名、SUBSCRIBE Topic Filter、PUBLISH Topic/Payload 全部错位。    FOR uiIndex := 0 TO uiStringLen - 1 DO        uiReadIndex := uiOffset + uiIndex;        IF uiReadIndex >= uiBufferLen THEN            F_MqttReadString := FALSE;            RETURN;        END_IF        sValue[uiIndex] := aBuffer[uiReadIndex];    END_FOR    sValue[uiStringLen] := 0;END_IFuiOffset := uiOffset + uiStringLen;F_MqttReadString := TRUE;

完整代码 15: F_MqttSkipVariableByteInteger.st

/// =======================================================================/// 名称      : F_MqttSkipVariableByteInteger/// 功能      : 跳过 MQTT 变长整数编码字段/// 说明      : MQTT 5.0 属性长度采用变长整数编码,当前轻量兼容层只需要校验并跳过该长度字段。/// 编程人员  : ControlRookie/// 时间      : 2026-05-08/// 版本      : V1.0/// ======================================================================={attribute 'hide_all_locals'}FUNCTION F_MqttSkipVariableByteInteger : BOOLVAR_INPUT    uiBufferLen  : UINT; // 当前 MQTT 完整报文长度或可用缓冲长度[byte]END_VARVAR_IN_OUT    aBuffer      : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区    uiOffset     : UINT; // 输入为变长整数起始偏移,成功后推进到变长整数之后[byte]END_VARVAR_OUTPUT    udiValue     : UDINT; // 解码出的变长整数数值,MQTT 5.0 属性长度使用该值[byte]END_VARVAR    udiMultiplier : UDINT; // MQTT 变长整数倍率,依次为 1、128、16384、2097152    byEncoded     : BYTE; // 当前读取的编码字节    uiLoop        : UINT; // 变长整数最多允许 4 个字节END_VAR// === IMPLEMENTATION ===udiValue := 0;udiMultiplier := 1;IF uiOffset >= uiBufferLen THEN    F_MqttSkipVariableByteInteger := FALSE;    RETURN;END_IFFOR uiLoop := 1 TO 4 DO    IF uiOffset >= uiBufferLen THEN        F_MqttSkipVariableByteInteger := FALSE;        RETURN;    END_IF    byEncoded := aBuffer[uiOffset];    udiValue := udiValue + TO_UDINT(byEncoded AND 16#7F) * udiMultiplier;    uiOffset := uiOffset + 1;    IF (byEncoded AND 16#80) = 0 THEN        F_MqttSkipVariableByteInteger := TRUE;        RETURN;    END_IF    udiMultiplier := udiMultiplier * 128;END_FORF_MqttSkipVariableByteInteger := FALSE;

完整代码 16: F_MqttStartsWith.st

/// =======================================================================/// 名称      : F_MqttStartsWith/// 功能      : 判断字符串是否以指定前缀开头/// 说明      : 避免依赖目标 IDE 的 LEFT 字符串库差异,供轻量 ACL 使用。/// 编程人员  : ControlRookie/// 时间      : 2026-05-08/// 版本      : V1.0/// ======================================================================={attribute 'hide_all_locals'}FUNCTION F_MqttStartsWith : BOOLVAR_INPUT    sValue  : STRING; // 需要检查的完整字符串    sPrefix : STRING; // 期望匹配的前缀,空前缀表示全部匹配END_VARVAR    uiValueLen  : UINT// 完整字符串长度[byte]    uiPrefixLen : UINT// 前缀字符串长度[byte]    uiIndex     : UINT// 字符逐字节比较索引,CODESYS STRING 字符下标从 0 开始[byte]END_VAR// === IMPLEMENTATION ===uiValueLen := TO_UINT(LEN(sValue));uiPrefixLen := TO_UINT(LEN(sPrefix));IF uiPrefixLen = 0 THEN    F_MqttStartsWith := TRUE;    RETURN;END_IFIF uiValueLen < uiPrefixLen THEN    F_MqttStartsWith := FALSE;    RETURN;END_IFFOR uiIndex := 0 TO uiPrefixLen - 1 DO    IF sValue[uiIndex] <> sPrefix[uiIndex] THEN        F_MqttStartsWith := FALSE;        RETURN;    END_IFEND_FORF_MqttStartsWith := TRUE;

这一篇你最该记住的几句话

  • Broker 源码不要按“文件夹顺序”读,要按“入口、状态、数据、报文、路由、事务”读。
  • CodeSys ST 工程最容易失控的不是语法,而是对象职责边界混乱。
  • 只要你能把本篇源码对象和在线变量对应起来,后续排查连接、订阅、发布和 QoS 问题就不会乱。

系列导航

  • 第 1 篇:源码加更01_Broker 工程入口、容量边界和数据模型
  • 第 2 篇:源码加更02_FB_MqttBroker 顶层调度、连接池和权限边界
  • 第 3 篇:源码加更03_单连接槽位、TCP 字节流和发送队列
  • 第 4 篇:源码加更04_MQTT 编解码器和字节工具函数
  • 第 5 篇:源码加更05_订阅表、Retain、PUBLISH 路由和业务事件
  • 第 6 篇:源码加更06_QoS 事务调度、重试和生产级闭环