这一组源码加更只有一个目标:把 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
iecst
/// =======================================================================
/// 名称 : M_BuildPublish
/// 功能 : 构建 Broker 出站 PUBLISH 报文
/// 说明 : 根据路由后的发布帧生成投递给订阅者的 MQTT PUBLISH。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-08
/// 版本 : V1.0
/// =======================================================================
{attribute 'hide_all_locals'}
METHOD M_BuildPublish : BOOL
VAR_INPUT
stPublish : ST_MqttBrokerPublishFrame; // 待投递给订阅者的发布帧
byProtocolLevel : BYTE; // 目标客户端 MQTT 协议级别,5 表示 PUBLISH 可变头需要追加零属性长度
uiWriteOffset : UINT; // 当前 PUBLISH 帧写入发送缓冲区的起始偏移,批量组包时用于追加到上一帧之后[byte]
udiBufferSize : UDINT; // 发送缓冲区总容量[byte]
END_VAR
VAR_IN_OUT
aBuffer : ARRAY[*] OF BYTE; // MQTT 发送缓冲区
END_VAR
VAR_OUTPUT
uiFrameLen : UINT; // 构建出的 PUBLISH 报文长度[byte]
END_VAR
VAR
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_IF
IF NOT F_MqttIsValidTopicName(sTopic := stPublish.sTopic) THEN
M_BuildPublish := FALSE;
RETURN;
END_IF
udiRemaining := 2 + TO_UDINT(stPublish.uiTopicLen) + TO_UDINT(stPublish.uiPayloadLen);
IF stPublish.eQoS <> E_MqttQoS.byQoS0 THEN
udiRemaining := udiRemaining + 2;
END_IF
IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THEN
// MQTT 5.0 PUBLISH 可变头在 Topic/PacketId 后必须携带 Properties。
// 当前 Broker 不发送任何 5.0 属性,因此属性长度固定编码为 0,占 1 字节。
udiRemaining := udiRemaining + 1;
END_IF
IF NOT F_MqttEncodeRemainingLength(
aBuffer := aRemaining,
udiValue := udiRemaining,
udiBufferSize := SIZEOF(aRemaining),
uiEncodedLen => uiRemainingLen) THEN
M_BuildPublish := FALSE;
RETURN;
END_IF
udiFrameLen := 1 + TO_UDINT(uiRemainingLen) + udiRemaining;
IF (TO_UDINT(uiWriteOffset) + udiFrameLen) > udiBufferSize THEN
M_BuildPublish := FALSE;
RETURN;
END_IF
aBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.byPublish);
IF stPublish.xDup THEN
aBuffer[uiWriteOffset] := aBuffer[uiWriteOffset] OR 16#08;
END_IF
CASE 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_CASE
IF stPublish.xRetain THEN
aBuffer[uiWriteOffset] := aBuffer[uiWriteOffset] OR 16#01;
END_IF
uiIndex := 0;
WHILE uiIndex < uiRemainingLen DO
aBuffer[uiWriteOffset + 1 + uiIndex] := aRemaining[uiIndex];
uiIndex := uiIndex + 1;
END_WHILE
uiOffset := uiWriteOffset + 1 + uiRemainingLen;
IF NOT F_MqttAppendString(
aBuffer := aBuffer,
uiOffset := uiOffset,
sValue := stPublish.sTopic,
udiBufferSize := udiBufferSize) THEN
M_BuildPublish := FALSE;
RETURN;
END_IF
IF 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_IF
IF 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_IF
IF 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_IF
uiFrameLen := uiOffset - uiWriteOffset;
uiLastFrameLen := uiFrameLen;
M_BuildPublish := TRUE;完整代码 02: FB_MqttBrokerCodec.M_BuildSimpleAck.st
iecst
/// =======================================================================
/// 名称 : M_BuildSimpleAck
/// 功能 : 构建固定长度 MQTT 协议响应
/// 说明 : 用于 CONNACK、PUBACK、SUBACK、UNSUBACK、PINGRESP 等轻量响应包。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-08
/// 版本 : V1.0
/// =======================================================================
{attribute 'hide_all_locals'}
METHOD M_BuildSimpleAck : BOOL
VAR_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_VAR
VAR_IN_OUT
aBuffer : ARRAY[*] OF BYTE; // MQTT 发送缓冲区
aReturnCodes : ARRAY[*] OF BYTE; // SUBACK 多 Topic 返回码数组,普通响应可传空闲数组
END_VAR
VAR_OUTPUT
uiFrameLen : UINT; // 构建出的 MQTT 响应报文长度[byte]
END_VAR
VAR
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_CASE
uiLastFrameLen := uiFrameLen;
M_BuildSimpleAck := TRUE;完整代码 03: FB_MqttBrokerCodec.M_ParseConnect.st
iecst
/// =======================================================================
/// 名称 : 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 : BOOL
VAR_INPUT
uiBodyOffset : UINT; // CONNECT 可变头在报文缓冲区中的起始偏移[byte]
uiFrameLen : UINT; // 当前 CONNECT 完整报文长度[byte]
END_VAR
VAR_IN_OUT
aBuffer : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区
stConnection : ST_MqttBrokerConnection; // 当前连接槽位,会被写入 ClientID、KeepAlive 和 Will 信息
END_VAR
VAR_OUTPUT
eError : E_MqttBrokerError; // 解析失败时返回的 Broker 错误码
END_VAR
VAR
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_IF
IF (uiOffset + 1) >= uiFrameLen THEN
eError := E_MqttBrokerError.uiProtocolMalformed;
M_ParseConnect := FALSE;
RETURN;
END_IF
uiProtocolLen := 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_CASE
IF uiOffset >= uiFrameLen THEN
eError := E_MqttBrokerError.uiProtocolMalformed;
M_ParseConnect := FALSE;
RETURN;
END_IF
byProtocolLevel := 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_IF
stConnection.byProtocolLevel := byProtocolLevel;
uiOffset := uiOffset + 1;
IF uiOffset >= uiFrameLen THEN
eError := E_MqttBrokerError.uiProtocolMalformed;
M_ParseConnect := FALSE;
RETURN;
END_IF
byFlags := aBuffer[uiOffset];
uiOffset := uiOffset + 1;
IF (byFlags AND 16#01) <> 0 THEN
eError := E_MqttBrokerError.uiProtocolMalformed;
M_ParseConnect := FALSE;
RETURN;
END_IF
IF (uiOffset + 1) >= uiFrameLen THEN
eError := E_MqttBrokerError.uiProtocolMalformed;
M_ParseConnect := FALSE;
RETURN;
END_IF
stConnection.uiKeepAlive := TO_UINT(aBuffer[uiOffset]) * 256 + TO_UINT(aBuffer[uiOffset + 1]);
IF stConnection.uiKeepAlive = 0 THEN
stConnection.uiKeepAlive := GVL_MqttBroker.cnDefaultKeepAlive;
END_IF
uiOffset := 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_IF
IF 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_IF
IF uiClientIdLen = 0 THEN
eError := E_MqttBrokerError.uiInvalidClientId;
M_ParseConnect := FALSE;
RETURN;
END_IF
stConnection.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#18, 3);
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_CASE
IF 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_IF
END_IF
IF (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_IF
IF (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_IF
uiLastFrameLen := uiFrameLen;
M_ParseConnect := TRUE;完整代码 04: FB_MqttBrokerCodec.M_ParsePublish.st
iecst
/// =======================================================================
/// 名称 : M_ParsePublish
/// 功能 : 解析 MQTT PUBLISH 报文
/// 说明 : 提取 Topic、Payload、QoS、Retain、DUP 和 PacketId,供路由层使用。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-08
/// 版本 : V1.0
/// =======================================================================
{attribute 'hide_all_locals'}
METHOD M_ParsePublish : BOOL
VAR_INPUT
uiSourceSlot : UINT; // 发布来源客户端槽位编号[1..cnMaxClientSlots]
byProtocolLevel : BYTE; // 当前连接协商的 MQTT 协议级别,5 表示需要跳过 PUBLISH Properties
uiBodyOffset : UINT; // PUBLISH 可变头起始偏移[byte]
uiFrameLen : UINT; // 当前 PUBLISH 完整报文长度[byte]
END_VAR
VAR_IN_OUT
aBuffer : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区
stPublish : ST_MqttBrokerPublishFrame; // 解析后的标准发布帧
END_VAR
VAR_OUTPUT
eError : E_MqttBrokerError; // 解析失败时返回的 Broker 错误码
END_VAR
VAR
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_IF
stPublish.xDup := (aBuffer[0] AND 16#08) <> 0;
stPublish.xRetain := (aBuffer[0] AND 16#01) <> 0;
byQoS := SHR(aBuffer[0] AND 16#06, 1);
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_CASE
IF 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_IF
IF NOT F_MqttIsValidTopicName(sTopic := stPublish.sTopic) THEN
eError := E_MqttBrokerError.uiInvalidTopic;
M_ParsePublish := FALSE;
RETURN;
END_IF
stPublish.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_IF
END_IF
IF 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_IF
stPublish.uiPayloadLen := uiFrameLen - uiOffset;
IF stPublish.uiPayloadLen > GVL_MqttBroker.cnMaxPayloadLen THEN
eError := E_MqttBrokerError.uiPacketTooLarge;
M_ParsePublish := FALSE;
RETURN;
END_IF
IF 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_IF
stPublish.xValid := TRUE;
uiLastFrameLen := uiFrameLen;
M_ParsePublish := TRUE;完整代码 05: FB_MqttBrokerCodec.M_ParseSubscribe.st
iecst
/// =======================================================================
/// 名称 : M_ParseSubscribe
/// 功能 : 解析 MQTT SUBSCRIBE 报文
/// 说明 : 第二阶段支持单个 SUBSCRIBE 报文内多个 Topic Filter,并逐项输出返回码。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-08
/// 版本 : V1.0
/// =======================================================================
{attribute 'hide_all_locals'}
METHOD M_ParseSubscribe : BOOL
VAR_INPUT
byProtocolLevel : BYTE; // 当前连接协商的 MQTT 协议级别,5 表示需要跳过 SUBSCRIBE Properties
uiBodyOffset : UINT; // SUBSCRIBE 可变头起始偏移[byte]
uiFrameLen : UINT; // 当前 SUBSCRIBE 完整报文长度[byte]
END_VAR
VAR_IN_OUT
aBuffer : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区
aTopicItems : ARRAY[*] OF ST_MqttBrokerTopicItem; // 解析出的多 Topic 订阅条目数组
END_VAR
VAR_OUTPUT
uiPacketId : UINT; // SUBSCRIBE Packet Identifier
uiItemCount : UINT; // 本次 SUBSCRIBE 成功解析出的 Topic 条目数量
eError : E_MqttBrokerError; // 解析失败时返回的 Broker 错误码
END_VAR
VAR
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_FOR
IF (uiOffset + 1) >= uiFrameLen THEN
eError := E_MqttBrokerError.uiProtocolMalformed;
M_ParseSubscribe := FALSE;
RETURN;
END_IF
uiPacketId := 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_IF
IF 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_IF
WHILE 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_WHILE
IF uiItemCount = 0 THEN
eError := E_MqttBrokerError.uiProtocolMalformed;
M_ParseSubscribe := FALSE;
RETURN;
END_IF
M_ParseSubscribe := TRUE;完整代码 06: FB_MqttBrokerCodec.M_ParseUnsubscribe.st
iecst
/// =======================================================================
/// 名称 : M_ParseUnsubscribe
/// 功能 : 解析 MQTT UNSUBSCRIBE 报文
/// 说明 : 第二阶段支持单个 UNSUBSCRIBE 报文内多个 Topic Filter。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-08
/// 版本 : V1.0
/// =======================================================================
{attribute 'hide_all_locals'}
METHOD M_ParseUnsubscribe : BOOL
VAR_INPUT
byProtocolLevel : BYTE; // 当前连接协商的 MQTT 协议级别,5 表示需要跳过 UNSUBSCRIBE Properties
uiBodyOffset : UINT; // UNSUBSCRIBE 可变头起始偏移[byte]
uiFrameLen : UINT; // 当前 UNSUBSCRIBE 完整报文长度[byte]
END_VAR
VAR_IN_OUT
aBuffer : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区
aTopicItems : ARRAY[*] OF ST_MqttBrokerTopicItem; // 解析出的多 Topic 取消订阅条目数组
END_VAR
VAR_OUTPUT
uiPacketId : UINT; // UNSUBSCRIBE Packet Identifier
uiItemCount : UINT; // 本次 UNSUBSCRIBE 成功解析出的 Topic 条目数量
eError : E_MqttBrokerError; // 解析失败时返回的 Broker 错误码
END_VAR
VAR
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_FOR
IF (uiOffset + 1) >= uiFrameLen THEN
eError := E_MqttBrokerError.uiProtocolMalformed;
M_ParseUnsubscribe := FALSE;
RETURN;
END_IF
uiPacketId := 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_IF
IF 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_IF
WHILE 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_WHILE
IF uiItemCount = 0 THEN
eError := E_MqttBrokerError.uiProtocolMalformed;
M_ParseUnsubscribe := FALSE;
RETURN;
END_IF
M_ParseUnsubscribe := TRUE;完整代码 07: FB_MqttBrokerCodec.st
iecst
/// =======================================================================
/// 名称 : FB_MqttBrokerCodec
/// 功能 : MQTT Broker 报文编解码器
/// 说明 : 本功能块不保存业务状态,只通过方法解析入站报文和构建服务端回包。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-08
/// 版本 : V1.0
/// =======================================================================
{attribute 'hide_all_locals'}
FUNCTION_BLOCK FB_MqttBrokerCodec
VAR
uiLastFrameLen : UINT; // 最近一次成功构建或解析的 MQTT 报文长度[byte]
END_VAR
// === IMPLEMENTATION ===完整代码 08: F_MqttAppendString.st
iecst
/// =======================================================================
/// 名称 : F_MqttAppendString
/// 功能 : 向 MQTT 报文缓冲区追加 UTF-8 字符串字段
/// 说明 : MQTT 字符串格式为 2 字节大端长度 + 字符串内容,本函数负责边界保护。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-08
/// 版本 : V1.0
/// =======================================================================
{attribute 'hide_all_locals'}
FUNCTION F_MqttAppendString : BOOL
VAR_INPUT
sValue : STRING; // 需要追加到 MQTT 报文中的字符串内容
udiBufferSize : UDINT; // 调用方报文缓冲区总容量[byte]
END_VAR
VAR_IN_OUT
aBuffer : ARRAY[*] OF BYTE; // 调用方报文缓冲区
uiOffset : UINT; // 当前写入偏移,成功后推进到字符串末尾后一字节[byte]
END_VAR
VAR
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_IF
aBuffer[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_FOR
END_IF
uiOffset := uiOffset + uiLen;
F_MqttAppendString := TRUE;完整代码 09: F_MqttContainsWildcard.st
iecst
/// =======================================================================
/// 名称 : F_MqttContainsWildcard
/// 功能 : 判断 Topic Filter 是否包含 MQTT 通配符
/// 说明 : 供 ACL 判断是否允许 + 或 # 通配符订阅。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-08
/// 版本 : V1.0
/// =======================================================================
{attribute 'hide_all_locals'}
FUNCTION F_MqttContainsWildcard : BOOL
VAR_INPUT
sTopicFilter : STRING; // 需要检查的 Topic Filter
END_VAR
VAR
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_IF
FOR uiIndex := 0 TO uiLen - 1 DO
IF (sTopicFilter[uiIndex] = 16#2B) OR (sTopicFilter[uiIndex] = 16#23) THEN
F_MqttContainsWildcard := TRUE;
RETURN;
END_IF
END_FOR
F_MqttContainsWildcard := FALSE;完整代码 10: F_MqttDecodeRemainingLength.st
iecst
/// =======================================================================
/// 名称 : F_MqttDecodeRemainingLength
/// 功能 : 解码 MQTT Remaining Length
/// 说明 : 从固定报头第 2 字节开始解析 MQTT 变长整数,最多读取 4 字节。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-08
/// 版本 : V1.0
/// =======================================================================
{attribute 'hide_all_locals'}
FUNCTION F_MqttDecodeRemainingLength : BOOL
VAR_INPUT
uiStartIndex : UINT; // Remaining Length 在缓冲区中的起始索引,通常为 1[byte]
uiBufferLen : UINT; // 当前缓冲区内有效数据长度[byte]
END_VAR
VAR_IN_OUT
aBuffer : ARRAY[*] OF BYTE; // MQTT 原始接收缓冲区
END_VAR
VAR_OUTPUT
udiValue : UDINT; // 解码后的 Remaining Length 数值[byte]
uiBytesUsed : UINT; // Remaining Length 字段实际占用的字节数[byte]
xNeedMore : BOOL; // 数据不足时置 TRUE,上层应继续读取 TCP 数据
END_VAR
VAR
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_IF
FOR 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_FOR
F_MqttDecodeRemainingLength := FALSE;完整代码 11: F_MqttEncodeRemainingLength.st
iecst
/// =======================================================================
/// 名称 : F_MqttEncodeRemainingLength
/// 功能 : 编码 MQTT Remaining Length
/// 说明 : 把 MQTT 剩余长度编码为 1~4 字节变长整数,并写入调用方提供的缓冲区。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-08
/// 版本 : V1.0
/// =======================================================================
{attribute 'hide_all_locals'}
FUNCTION F_MqttEncodeRemainingLength : BOOL
VAR_INPUT
udiValue : UDINT; // 需要编码的 MQTT Remaining Length 数值[byte]
udiBufferSize : UDINT; // 调用方输出缓冲区容量[byte]
END_VAR
VAR_IN_OUT
aBuffer : ARRAY[*] OF BYTE; // 调用方输出缓冲区,函数从索引 0 起写入编码结果
END_VAR
VAR_OUTPUT
uiEncodedLen : UINT; // 实际编码产生的字节数[byte]
END_VAR
VAR
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_IF
REPEAT
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 = 0
END_REPEAT
uiEncodedLen := uiIndex;
F_MqttEncodeRemainingLength := TRUE;完整代码 12: F_MqttIsValidTopicFilter.st
iecst
/// =======================================================================
/// 名称 : F_MqttIsValidTopicFilter
/// 功能 : 校验 MQTT Topic Filter
/// 说明 : SUBSCRIBE 使用的过滤器允许 + / #,但必须满足 MQTT 通配符位置规则。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-08
/// 版本 : V1.0
/// =======================================================================
{attribute 'hide_all_locals'}
FUNCTION F_MqttIsValidTopicFilter : BOOL
VAR_INPUT
sFilter : STRING; // 待校验的 MQTT Topic Filter
END_VAR
VAR
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_IF
IF uiLen > GVL_MqttBroker.cnMaxTopicLen THEN
F_MqttIsValidTopicFilter := FALSE;
RETURN;
END_IF
FOR 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_CASE
END_FOR
F_MqttIsValidTopicFilter := TRUE;完整代码 13: F_MqttIsValidTopicName.st
iecst
/// =======================================================================
/// 名称 : F_MqttIsValidTopicName
/// 功能 : 校验 MQTT Topic Name
/// 说明 : PUBLISH 使用的 Topic Name 不能为空,且不能包含 + / # 通配符。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-08
/// 版本 : V1.0
/// =======================================================================
{attribute 'hide_all_locals'}
FUNCTION F_MqttIsValidTopicName : BOOL
VAR_INPUT
sTopic : STRING; // 待校验的 MQTT Topic Name
END_VAR
VAR
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_IF
IF uiLen > GVL_MqttBroker.cnMaxTopicLen THEN
F_MqttIsValidTopicName := FALSE;
RETURN;
END_IF
FOR uiIndex := 0 TO uiLen - 1 DO
IF (sTopic[uiIndex] = 16#2B) OR (sTopic[uiIndex] = 16#23) THEN
F_MqttIsValidTopicName := FALSE;
RETURN;
END_IF
END_FOR
F_MqttIsValidTopicName := TRUE;完整代码 14: F_MqttReadString.st
iecst
/// =======================================================================
/// 名称 : F_MqttReadString
/// 功能 : 从 MQTT 报文缓冲区读取 UTF-8 字符串字段
/// 说明 : 读取 2 字节大端长度和后续内容,并推进调用方偏移。
/// 注意 : CodeSys/CODESYS STRING 字符下标按 0 基访问,不能按 1 基复制。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-08
/// 版本 : V1.0
/// =======================================================================
{attribute 'hide_all_locals'}
FUNCTION F_MqttReadString : BOOL
VAR_INPUT
uiBufferLen : UINT; // 当前 MQTT 报文有效长度[byte]
uiMaxLen : UINT; // 输出字符串允许保存的最大长度[byte]
END_VAR
VAR_IN_OUT
aBuffer : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区
uiOffset : UINT; // 当前读取偏移,成功后推进到字符串末尾后一字节[byte]
sValue : STRING; // 读取出的字符串内容,短字段由调用方使用默认长度临时缓冲转接
END_VAR
VAR_OUTPUT
uiStringLen : UINT; // MQTT 字符串字段声明的原始长度[byte]
END_VAR
VAR
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_IF
uiStringLen := TO_UINT(aBuffer[uiOffset]) * 256 + TO_UINT(aBuffer[uiOffset + 1]);
uiOffset := uiOffset + 2;
IF uiStringLen > uiMaxLen THEN
F_MqttReadString := FALSE;
RETURN;
END_IF
IF (TO_UDINT(uiOffset) + TO_UDINT(uiStringLen)) > TO_UDINT(uiBufferLen) THEN
F_MqttReadString := FALSE;
RETURN;
END_IF
IF 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_IF
uiOffset := uiOffset + uiStringLen;
F_MqttReadString := TRUE;完整代码 15: F_MqttSkipVariableByteInteger.st
iecst
/// =======================================================================
/// 名称 : F_MqttSkipVariableByteInteger
/// 功能 : 跳过 MQTT 变长整数编码字段
/// 说明 : MQTT 5.0 属性长度采用变长整数编码,当前轻量兼容层只需要校验并跳过该长度字段。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-08
/// 版本 : V1.0
/// =======================================================================
{attribute 'hide_all_locals'}
FUNCTION F_MqttSkipVariableByteInteger : BOOL
VAR_INPUT
uiBufferLen : UINT; // 当前 MQTT 完整报文长度或可用缓冲长度[byte]
END_VAR
VAR_IN_OUT
aBuffer : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区
uiOffset : UINT; // 输入为变长整数起始偏移,成功后推进到变长整数之后[byte]
END_VAR
VAR_OUTPUT
udiValue : UDINT; // 解码出的变长整数数值,MQTT 5.0 属性长度使用该值[byte]
END_VAR
VAR
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_IF
FOR 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_FOR
F_MqttSkipVariableByteInteger := FALSE;完整代码 16: F_MqttStartsWith.st
iecst
/// =======================================================================
/// 名称 : F_MqttStartsWith
/// 功能 : 判断字符串是否以指定前缀开头
/// 说明 : 避免依赖目标 IDE 的 LEFT 字符串库差异,供轻量 ACL 使用。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-08
/// 版本 : V1.0
/// =======================================================================
{attribute 'hide_all_locals'}
FUNCTION F_MqttStartsWith : BOOL
VAR_INPUT
sValue : STRING; // 需要检查的完整字符串
sPrefix : STRING; // 期望匹配的前缀,空前缀表示全部匹配
END_VAR
VAR
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_IF
IF uiValueLen < uiPrefixLen THEN
F_MqttStartsWith := FALSE;
RETURN;
END_IF
FOR uiIndex := 0 TO uiPrefixLen - 1 DO
IF sValue[uiIndex] <> sPrefix[uiIndex] THEN
F_MqttStartsWith := FALSE;
RETURN;
END_IF
END_FOR
F_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 事务调度、重试和生产级闭环
评论区预留
这里先保留评论和回复结构,不接入第三方服务。后续统一决定登录、匿名、审核、反垃圾和静态站兼容策略。