源码加更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_MqttBrokerCodec.M_BuildPublish.st | ||
FB_MqttBrokerCodec.M_BuildSimpleAck.st | ||
FB_MqttBrokerCodec.M_ParseConnect.st | ||
FB_MqttBrokerCodec.M_ParsePublish.st | ||
FB_MqttBrokerCodec.M_ParseSubscribe.st | ||
FB_MqttBrokerCodec.M_ParseUnsubscribe.st | ||
FB_MqttBrokerCodec.st | ||
F_MqttAppendString.st | ||
F_MqttContainsWildcard.st | ||
F_MqttDecodeRemainingLength.st | ||
F_MqttEncodeRemainingLength.st | ||
F_MqttIsValidTopicFilter.st | ||
F_MqttIsValidTopicName.st | ||
F_MqttReadString.st | ||
F_MqttSkipVariableByteInteger.st | ||
F_MqttStartsWith.st |
推荐阅读顺序
先看 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_INPUTstPublish : ST_MqttBrokerPublishFrame; // 待投递给订阅者的发布帧byProtocolLevel : BYTE; // 目标客户端 MQTT 协议级别,5 表示 PUBLISH 可变头需要追加零属性长度uiWriteOffset : UINT; // 当前 PUBLISH 帧写入发送缓冲区的起始偏移,批量组包时用于追加到上一帧之后[byte]udiBufferSize : UDINT; // 发送缓冲区总容量[byte]END_VARVAR_IN_OUTaBuffer : ARRAY[*] OF BYTE; // MQTT 发送缓冲区END_VARVAR_OUTPUTuiFrameLen : UINT; // 构建出的 PUBLISH 报文长度[byte]END_VARVARaRemaining : 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 THENM_BuildPublish := FALSE;RETURN;END_IFIF NOT F_MqttIsValidTopicName(sTopic := stPublish.sTopic) THENM_BuildPublish := FALSE;RETURN;END_IFudiRemaining := 2 + TO_UDINT(stPublish.uiTopicLen) + TO_UDINT(stPublish.uiPayloadLen);IF stPublish.eQoS <> E_MqttQoS.byQoS0 THENudiRemaining := 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) THENM_BuildPublish := FALSE;RETURN;END_IFudiFrameLen := 1 + TO_UDINT(uiRemainingLen) + udiRemaining;IF (TO_UDINT(uiWriteOffset) + udiFrameLen) > udiBufferSize THENM_BuildPublish := FALSE;RETURN;END_IFaBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.byPublish);IF stPublish.xDup THENaBuffer[uiWriteOffset] := aBuffer[uiWriteOffset] OR 16#08;END_IFCASE stPublish.eQoS OFE_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;ELSEM_BuildPublish := FALSE;RETURN;END_CASEIF stPublish.xRetain THENaBuffer[uiWriteOffset] := aBuffer[uiWriteOffset] OR 16#01;END_IFuiIndex := 0;WHILE uiIndex < uiRemainingLen DOaBuffer[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) THENM_BuildPublish := FALSE;RETURN;END_IFIF stPublish.eQoS <> E_MqttQoS.byQoS0 THENIF (TO_UDINT(uiOffset) + 2) > udiBufferSize THENM_BuildPublish := FALSE;RETURN;END_IFaBuffer[uiOffset] := TO_BYTE(stPublish.uiPacketId / 256);aBuffer[uiOffset + 1] := TO_BYTE(stPublish.uiPacketId MOD 256);uiOffset := uiOffset + 2;END_IFIF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENIF (TO_UDINT(uiOffset) + 1) > udiBufferSize THENM_BuildPublish := FALSE;RETURN;END_IFaBuffer[uiOffset] := 0;uiOffset := uiOffset + 1;END_IFIF stPublish.uiPayloadLen > 0 THENIF (TO_UDINT(uiOffset) + TO_UDINT(stPublish.uiPayloadLen)) > udiBufferSize THENM_BuildPublish := FALSE;RETURN;END_IF// Payload 存在 ST STRING 中时按 0 基下标读取。// 这和 MQTT 报文字节数组 aBuffer[0..] 的下标体系一致,可以避免转发时首字节丢失、尾部多出垃圾字符。FOR uiIndex := 0 TO stPublish.uiPayloadLen - 1 DOaBuffer[uiOffset + uiIndex] := stPublish.sPayload[uiIndex];END_FORuiOffset := 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_INPUTePacketType : E_MqttPacketType; // 需要构建的 MQTT 响应报文类型byProtocolLevel : BYTE; // 目标客户端 MQTT 协议级别,5 表示响应中需要携带零属性长度uiPacketId : UINT; // 需要带 Packet Identifier 的响应使用该值,无需 PacketId 时为 0byReturnCode : BYTE; // CONNACK/SUBACK 返回码,其他响应通常为 0uiReturnCount : UINT; // SUBACK 多 Topic 返回码数量,普通响应传 0uiWriteOffset : UINT; // 当前响应帧写入发送缓冲区的起始偏移,批量组包时用于追加到上一帧之后[byte]udiBufferSize : UDINT; // 发送缓冲区总容量[byte]END_VARVAR_IN_OUTaBuffer : ARRAY[*] OF BYTE; // MQTT 发送缓冲区aReturnCodes : ARRAY[*] OF BYTE; // SUBACK 多 Topic 返回码数组,普通响应可传空闲数组END_VARVAR_OUTPUTuiFrameLen : UINT; // 构建出的 MQTT 响应报文长度[byte]END_VARVARuiIndex : UINT; // SUBACK 多返回码复制索引[1..cnMaxTopicItemsPerPacket]udiNeededLen : UDINT; // 当前响应帧需要的完整缓冲长度,包含固定头、可变头和返回码[byte]END_VAR// === IMPLEMENTATION ===uiFrameLen := 0;CASE ePacketType OFE_MqttPacketType.byConnAck:IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENudiNeededLen := 5;ELSEudiNeededLen := 4;END_IFIF (TO_UDINT(uiWriteOffset) + udiNeededLen) > udiBufferSize THENM_BuildSimpleAck := FALSE;RETURN;END_IFaBuffer[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;ELSEaBuffer[uiWriteOffset + 1] := 2;uiFrameLen := 4;END_IFE_MqttPacketType.byPubAck:IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENudiNeededLen := 6;ELSEudiNeededLen := 4;END_IFIF (TO_UDINT(uiWriteOffset) + udiNeededLen) > udiBufferSize THENM_BuildSimpleAck := FALSE;RETURN;END_IFaBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.byPubAck);IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENaBuffer[uiWriteOffset + 1] := 4;ELSEaBuffer[uiWriteOffset + 1] := 2;END_IFaBuffer[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;ELSEuiFrameLen := 4;END_IFE_MqttPacketType.byPubRec:IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENudiNeededLen := 6;ELSEudiNeededLen := 4;END_IFIF (TO_UDINT(uiWriteOffset) + udiNeededLen) > udiBufferSize THENM_BuildSimpleAck := FALSE;RETURN;END_IFaBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.byPubRec);IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENaBuffer[uiWriteOffset + 1] := 4;ELSEaBuffer[uiWriteOffset + 1] := 2;END_IFaBuffer[uiWriteOffset + 2] := TO_BYTE(uiPacketId / 256);aBuffer[uiWriteOffset + 3] := TO_BYTE(uiPacketId MOD 256);IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENaBuffer[uiWriteOffset + 4] := byReturnCode;aBuffer[uiWriteOffset + 5] := 0;uiFrameLen := 6;ELSEuiFrameLen := 4;END_IFE_MqttPacketType.byPubRel:IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENudiNeededLen := 6;ELSEudiNeededLen := 4;END_IFIF (TO_UDINT(uiWriteOffset) + udiNeededLen) > udiBufferSize THENM_BuildSimpleAck := FALSE;RETURN;END_IFaBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.byPubRel);IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENaBuffer[uiWriteOffset + 1] := 4;ELSEaBuffer[uiWriteOffset + 1] := 2;END_IFaBuffer[uiWriteOffset + 2] := TO_BYTE(uiPacketId / 256);aBuffer[uiWriteOffset + 3] := TO_BYTE(uiPacketId MOD 256);IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENaBuffer[uiWriteOffset + 4] := byReturnCode;aBuffer[uiWriteOffset + 5] := 0;uiFrameLen := 6;ELSEuiFrameLen := 4;END_IFE_MqttPacketType.byPubComp:IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENudiNeededLen := 6;ELSEudiNeededLen := 4;END_IFIF (TO_UDINT(uiWriteOffset) + udiNeededLen) > udiBufferSize THENM_BuildSimpleAck := FALSE;RETURN;END_IFaBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.byPubComp);IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENaBuffer[uiWriteOffset + 1] := 4;ELSEaBuffer[uiWriteOffset + 1] := 2;END_IFaBuffer[uiWriteOffset + 2] := TO_BYTE(uiPacketId / 256);aBuffer[uiWriteOffset + 3] := TO_BYTE(uiPacketId MOD 256);IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENaBuffer[uiWriteOffset + 4] := byReturnCode;aBuffer[uiWriteOffset + 5] := 0;uiFrameLen := 6;ELSEuiFrameLen := 4;END_IFE_MqttPacketType.bySubAck:IF uiReturnCount = 0 THENIF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENudiNeededLen := 6;ELSEudiNeededLen := 5;END_IFELSEIF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENudiNeededLen := 5 + TO_UDINT(uiReturnCount);ELSEudiNeededLen := 4 + TO_UDINT(uiReturnCount);END_IFEND_IFIF (TO_UDINT(uiWriteOffset) + udiNeededLen) > udiBufferSize THENM_BuildSimpleAck := FALSE;RETURN;END_IFaBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.bySubAck);IF uiReturnCount = 0 THENIF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENaBuffer[uiWriteOffset + 1] := 4;ELSEaBuffer[uiWriteOffset + 1] := 3;END_IFELSEIF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENaBuffer[uiWriteOffset + 1] := TO_BYTE(3 + uiReturnCount);ELSEaBuffer[uiWriteOffset + 1] := TO_BYTE(2 + uiReturnCount);END_IFEND_IFaBuffer[uiWriteOffset + 2] := TO_BYTE(uiPacketId / 256);aBuffer[uiWriteOffset + 3] := TO_BYTE(uiPacketId MOD 256);IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENaBuffer[uiWriteOffset + 4] := 0;IF uiReturnCount = 0 THENaBuffer[uiWriteOffset + 5] := byReturnCode;uiFrameLen := 6;ELSEFOR uiIndex := 1 TO uiReturnCount DOIF uiIndex > GVL_MqttBroker.cnMaxTopicItemsPerPacket THENM_BuildSimpleAck := FALSE;RETURN;END_IFaBuffer[uiWriteOffset + 4 + uiIndex] := aReturnCodes[uiIndex];END_FORuiFrameLen := 5 + uiReturnCount;END_IFELSEIF uiReturnCount = 0 THENaBuffer[uiWriteOffset + 4] := byReturnCode;uiFrameLen := 5;ELSEFOR uiIndex := 1 TO uiReturnCount DOIF uiIndex > GVL_MqttBroker.cnMaxTopicItemsPerPacket THENM_BuildSimpleAck := FALSE;RETURN;END_IFaBuffer[uiWriteOffset + 3 + uiIndex] := aReturnCodes[uiIndex];END_FORuiFrameLen := 4 + uiReturnCount;END_IFEND_IFE_MqttPacketType.byUnsubAck:IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENudiNeededLen := 6;ELSEudiNeededLen := 4;END_IFIF (TO_UDINT(uiWriteOffset) + udiNeededLen) > udiBufferSize THENM_BuildSimpleAck := FALSE;RETURN;END_IFaBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.byUnsubAck);IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENaBuffer[uiWriteOffset + 1] := 4;ELSEaBuffer[uiWriteOffset + 1] := 2;END_IFaBuffer[uiWriteOffset + 2] := TO_BYTE(uiPacketId / 256);aBuffer[uiWriteOffset + 3] := TO_BYTE(uiPacketId MOD 256);IF byProtocolLevel = GVL_MqttBroker.cnMqttProtocolLevel5 THENaBuffer[uiWriteOffset + 4] := byReturnCode;aBuffer[uiWriteOffset + 5] := 0;uiFrameLen := 6;ELSEuiFrameLen := 4;END_IFE_MqttPacketType.byPingResp:IF (TO_UDINT(uiWriteOffset) + 2) > udiBufferSize THENM_BuildSimpleAck := FALSE;RETURN;END_IFaBuffer[uiWriteOffset] := TO_BYTE(E_MqttPacketType.byPingResp);aBuffer[uiWriteOffset + 1] := 0;uiFrameLen := 2;ELSEM_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_INPUTuiBodyOffset : UINT; // CONNECT 可变头在报文缓冲区中的起始偏移[byte]uiFrameLen : UINT; // 当前 CONNECT 完整报文长度[byte]END_VARVAR_IN_OUTaBuffer : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区stConnection : ST_MqttBrokerConnection; // 当前连接槽位,会被写入 ClientID、KeepAlive 和 Will 信息END_VARVAR_OUTPUTeError : E_MqttBrokerError; // 解析失败时返回的 Broker 错误码END_VARVARuiOffset : 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.0byWillQoS : BYTE; // 从 CONNECT Flags 中提取的 Will QoS 数值sTempString : STRING; // CONNECT 用户名/密码临时缓冲,再按目标字段容量赋值END_VAR// === IMPLEMENTATION ===eError := E_MqttBrokerError.uiNoError;uiOffset := uiBodyOffset;IF uiFrameLen < 14 THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseConnect := FALSE;RETURN;END_IFIF (uiOffset + 1) >= uiFrameLen THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseConnect := FALSE;RETURN;END_IFuiProtocolLen := TO_UINT(aBuffer[uiOffset]) * 256 + TO_UINT(aBuffer[uiOffset + 1]);CASE uiProtocolLen OF4:IF (uiOffset + 5) >= uiFrameLen THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseConnect := FALSE;RETURN;END_IFIF (aBuffer[uiOffset + 2] <> 16#4D)OR (aBuffer[uiOffset + 3] <> 16#51)OR (aBuffer[uiOffset + 4] <> 16#54)OR (aBuffer[uiOffset + 5] <> 16#54) THENeError := E_MqttBrokerError.uiUnsupportedProtocol;M_ParseConnect := FALSE;RETURN;END_IFuiOffset := uiOffset + 6;6:IF (uiOffset + 7) >= uiFrameLen THENeError := 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) THENeError := E_MqttBrokerError.uiUnsupportedProtocol;M_ParseConnect := FALSE;RETURN;END_IFuiOffset := uiOffset + 8;ELSEeError := E_MqttBrokerError.uiUnsupportedProtocol;M_ParseConnect := FALSE;RETURN;END_CASEIF uiOffset >= uiFrameLen THENeError := 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)) THENeError := E_MqttBrokerError.uiUnsupportedProtocol;M_ParseConnect := FALSE;RETURN;END_IFstConnection.byProtocolLevel := byProtocolLevel;uiOffset := uiOffset + 1;IF uiOffset >= uiFrameLen THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseConnect := FALSE;RETURN;END_IFbyFlags := aBuffer[uiOffset];uiOffset := uiOffset + 1;IF (byFlags AND 16#01) <> 0 THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseConnect := FALSE;RETURN;END_IFIF (uiOffset + 1) >= uiFrameLen THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseConnect := FALSE;RETURN;END_IFstConnection.uiKeepAlive := TO_UINT(aBuffer[uiOffset]) * 256 + TO_UINT(aBuffer[uiOffset + 1]);IF stConnection.uiKeepAlive = 0 THENstConnection.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) THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseConnect := FALSE;RETURN;END_IFIF (TO_UDINT(uiOffset) + udiPropertyLen) > TO_UDINT(uiFrameLen) THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseConnect := FALSE;RETURN;END_IFuiOffset := uiOffset + TO_UINT(udiPropertyLen);END_IFIF NOT F_MqttReadString(aBuffer := aBuffer,uiOffset := uiOffset,sValue := stConnection.sClientId,uiBufferLen := uiFrameLen,uiMaxLen := GVL_MqttBroker.cnMaxClientIdLen,uiStringLen => uiClientIdLen) THENeError := E_MqttBrokerError.uiInvalidClientId;M_ParseConnect := FALSE;RETURN;END_IFIF uiClientIdLen = 0 THENeError := 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#18, 3);CASE byWillQoS OF0:stConnection.eWillQoS := E_MqttQoS.byQoS0;1:stConnection.eWillQoS := E_MqttQoS.byQoS1;2:stConnection.eWillQoS := E_MqttQoS.byQoS2;ELSEeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseConnect := FALSE;RETURN;END_CASEIF stConnection.xWillFlag THENIF 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) THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseConnect := FALSE;RETURN;END_IFIF (TO_UDINT(uiOffset) + udiPropertyLen) > TO_UDINT(uiFrameLen) THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseConnect := FALSE;RETURN;END_IFuiOffset := uiOffset + TO_UINT(udiPropertyLen);END_IFIF NOT F_MqttReadString(aBuffer := aBuffer,uiOffset := uiOffset,sValue := stConnection.sWillTopic,uiBufferLen := uiFrameLen,uiMaxLen := GVL_MqttBroker.cnMaxTopicLen,uiStringLen => uiWillTopicLen) THENeError := E_MqttBrokerError.uiInvalidTopic;M_ParseConnect := FALSE;RETURN;END_IFIF NOT F_MqttIsValidTopicName(sTopic := stConnection.sWillTopic) THENeError := E_MqttBrokerError.uiInvalidTopic;M_ParseConnect := FALSE;RETURN;END_IFIF NOT F_MqttReadString(aBuffer := aBuffer,uiOffset := uiOffset,sValue := stConnection.sWillPayload,uiBufferLen := uiFrameLen,uiMaxLen := GVL_MqttBroker.cnMaxPayloadLen,uiStringLen => uiWillMsgLen) THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseConnect := FALSE;RETURN;END_IFEND_IFIF (byFlags AND 16#80) <> 0 THENIF NOT F_MqttReadString(aBuffer := aBuffer,uiOffset := uiOffset,sValue := sTempString,uiBufferLen := uiFrameLen,uiMaxLen := GVL_MqttBroker.cnMaxUsernameLen,uiStringLen => uiUsernameLen) THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseConnect := FALSE;RETURN;END_IFstConnection.sUsername := sTempString;END_IFIF (byFlags AND 16#40) <> 0 THENIF NOT F_MqttReadString(aBuffer := aBuffer,uiOffset := uiOffset,sValue := sTempString,uiBufferLen := uiFrameLen,uiMaxLen := GVL_MqttBroker.cnMaxPasswordLen,uiStringLen => uiPasswordLen) THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseConnect := FALSE;RETURN;END_IFstConnection.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_INPUTuiSourceSlot : UINT; // 发布来源客户端槽位编号[1..cnMaxClientSlots]byProtocolLevel : BYTE; // 当前连接协商的 MQTT 协议级别,5 表示需要跳过 PUBLISH PropertiesuiBodyOffset : UINT; // PUBLISH 可变头起始偏移[byte]uiFrameLen : UINT; // 当前 PUBLISH 完整报文长度[byte]END_VARVAR_IN_OUTaBuffer : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区stPublish : ST_MqttBrokerPublishFrame; // 解析后的标准发布帧END_VARVAR_OUTPUTeError : E_MqttBrokerError; // 解析失败时返回的 Broker 错误码END_VARVARuiOffset : 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 THENeError := 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#06, 1);CASE byQoS OF0:stPublish.eQoS := E_MqttQoS.byQoS0;1:stPublish.eQoS := E_MqttQoS.byQoS1;2:stPublish.eQoS := E_MqttQoS.byQoS2;ELSEeError := 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) THENeError := E_MqttBrokerError.uiInvalidTopic;M_ParsePublish := FALSE;RETURN;END_IFIF NOT F_MqttIsValidTopicName(sTopic := stPublish.sTopic) THENeError := E_MqttBrokerError.uiInvalidTopic;M_ParsePublish := FALSE;RETURN;END_IFstPublish.uiTopicLen := uiTopicLen;IF stPublish.eQoS <> E_MqttQoS.byQoS0 THENIF (uiOffset + 1) >= uiFrameLen THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParsePublish := FALSE;RETURN;END_IFstPublish.uiPacketId := TO_UINT(aBuffer[uiOffset]) * 256 + TO_UINT(aBuffer[uiOffset + 1]);uiOffset := uiOffset + 2;IF stPublish.uiPacketId = 0 THENeError := 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) THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParsePublish := FALSE;RETURN;END_IFIF (TO_UDINT(uiOffset) + udiPropertyLen) > TO_UDINT(uiFrameLen) THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParsePublish := FALSE;RETURN;END_IFuiOffset := uiOffset + TO_UINT(udiPropertyLen);END_IFstPublish.uiPayloadLen := uiFrameLen - uiOffset;IF stPublish.uiPayloadLen > GVL_MqttBroker.cnMaxPayloadLen THENeError := 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 DOstPublish.sPayload[uiPayloadIdx] := aBuffer[uiOffset + uiPayloadIdx];END_FORstPublish.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_INPUTbyProtocolLevel : BYTE; // 当前连接协商的 MQTT 协议级别,5 表示需要跳过 SUBSCRIBE PropertiesuiBodyOffset : UINT; // SUBSCRIBE 可变头起始偏移[byte]uiFrameLen : UINT; // 当前 SUBSCRIBE 完整报文长度[byte]END_VARVAR_IN_OUTaBuffer : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区aTopicItems : ARRAY[*] OF ST_MqttBrokerTopicItem; // 解析出的多 Topic 订阅条目数组END_VARVAR_OUTPUTuiPacketId : UINT; // SUBSCRIBE Packet IdentifieruiItemCount : UINT; // 本次 SUBSCRIBE 成功解析出的 Topic 条目数量eError : E_MqttBrokerError; // 解析失败时返回的 Broker 错误码END_VARVARuiOffset : 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 DOaTopicItems[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 THENeError := 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 THENeError := 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) THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseSubscribe := FALSE;RETURN;END_IFIF (TO_UDINT(uiOffset) + udiPropertyLen) > TO_UDINT(uiFrameLen) THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseSubscribe := FALSE;RETURN;END_IFuiOffset := uiOffset + TO_UINT(udiPropertyLen);END_IFWHILE uiOffset < uiFrameLen DOIF uiItemCount >= GVL_MqttBroker.cnMaxTopicItemsPerPacket THENeError := E_MqttBrokerError.uiPacketTooLarge;M_ParseSubscribe := FALSE;RETURN;END_IFuiIndex := uiItemCount + 1;IF NOT F_MqttReadString(aBuffer := aBuffer,uiOffset := uiOffset,sValue := aTopicItems[uiIndex].sTopicFilter,uiBufferLen := uiFrameLen,uiMaxLen := GVL_MqttBroker.cnMaxTopicLen,uiStringLen => uiFilterLen) THENeError := E_MqttBrokerError.uiInvalidTopic;M_ParseSubscribe := FALSE;RETURN;END_IFIF uiOffset >= uiFrameLen THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseSubscribe := FALSE;RETURN;END_IFbyQoS := aBuffer[uiOffset];uiOffset := uiOffset + 1;CASE byQoS OF0:aTopicItems[uiIndex].eQoS := E_MqttQoS.byQoS0;1:aTopicItems[uiIndex].eQoS := E_MqttQoS.byQoS1;2:aTopicItems[uiIndex].eQoS := E_MqttQoS.byQoS2;ELSEeError := E_MqttBrokerError.uiUnsupportedQoS;M_ParseSubscribe := FALSE;RETURN;END_CASEIF NOT F_MqttIsValidTopicFilter(sFilter := aTopicItems[uiIndex].sTopicFilter) THENeError := E_MqttBrokerError.uiInvalidTopic;M_ParseSubscribe := FALSE;RETURN;END_IFaTopicItems[uiIndex].xUsed := TRUE;aTopicItems[uiIndex].uiFilterLen := uiFilterLen;aTopicItems[uiIndex].byReturnCode := TO_BYTE(aTopicItems[uiIndex].eQoS);uiItemCount := uiItemCount + 1;END_WHILEIF uiItemCount = 0 THENeError := 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_INPUTbyProtocolLevel : BYTE; // 当前连接协商的 MQTT 协议级别,5 表示需要跳过 UNSUBSCRIBE PropertiesuiBodyOffset : UINT; // UNSUBSCRIBE 可变头起始偏移[byte]uiFrameLen : UINT; // 当前 UNSUBSCRIBE 完整报文长度[byte]END_VARVAR_IN_OUTaBuffer : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区aTopicItems : ARRAY[*] OF ST_MqttBrokerTopicItem; // 解析出的多 Topic 取消订阅条目数组END_VARVAR_OUTPUTuiPacketId : UINT; // UNSUBSCRIBE Packet IdentifieruiItemCount : UINT; // 本次 UNSUBSCRIBE 成功解析出的 Topic 条目数量eError : E_MqttBrokerError; // 解析失败时返回的 Broker 错误码END_VARVARuiOffset : 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 DOaTopicItems[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 THENeError := 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 THENeError := 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) THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseUnsubscribe := FALSE;RETURN;END_IFIF (TO_UDINT(uiOffset) + udiPropertyLen) > TO_UDINT(uiFrameLen) THENeError := E_MqttBrokerError.uiProtocolMalformed;M_ParseUnsubscribe := FALSE;RETURN;END_IFuiOffset := uiOffset + TO_UINT(udiPropertyLen);END_IFWHILE uiOffset < uiFrameLen DOIF uiItemCount >= GVL_MqttBroker.cnMaxTopicItemsPerPacket THENeError := E_MqttBrokerError.uiPacketTooLarge;M_ParseUnsubscribe := FALSE;RETURN;END_IFuiIndex := uiItemCount + 1;IF NOT F_MqttReadString(aBuffer := aBuffer,uiOffset := uiOffset,sValue := aTopicItems[uiIndex].sTopicFilter,uiBufferLen := uiFrameLen,uiMaxLen := GVL_MqttBroker.cnMaxTopicLen,uiStringLen => uiFilterLen) THENeError := E_MqttBrokerError.uiInvalidTopic;M_ParseUnsubscribe := FALSE;RETURN;END_IFIF NOT F_MqttIsValidTopicFilter(sFilter := aTopicItems[uiIndex].sTopicFilter) THENeError := E_MqttBrokerError.uiInvalidTopic;M_ParseUnsubscribe := FALSE;RETURN;END_IFaTopicItems[uiIndex].xUsed := TRUE;aTopicItems[uiIndex].uiFilterLen := uiFilterLen;aTopicItems[uiIndex].eQoS := E_MqttQoS.byQoS0;aTopicItems[uiIndex].byReturnCode := 0;uiItemCount := uiItemCount + 1;END_WHILEIF uiItemCount = 0 THENeError := 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_MqttBrokerCodecVARuiLastFrameLen : 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_INPUTsValue : STRING; // 需要追加到 MQTT 报文中的字符串内容udiBufferSize : UDINT; // 调用方报文缓冲区总容量[byte]END_VARVAR_IN_OUTaBuffer : ARRAY[*] OF BYTE; // 调用方报文缓冲区uiOffset : UINT; // 当前写入偏移,成功后推进到字符串末尾后一字节[byte]END_VARVARuiLen : 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 THENF_MqttAppendString := FALSE;RETURN;END_IFaBuffer[uiOffset] := TO_BYTE(uiLen / 256);aBuffer[uiOffset + 1] := TO_BYTE(uiLen MOD 256);uiOffset := uiOffset + 2;IF uiLen > 0 THENFOR uiIndex := 0 TO uiLen - 1 DOuiWriteIndex := uiOffset + uiIndex;IF TO_UDINT(uiWriteIndex) >= udiBufferSize THENF_MqttAppendString := FALSE;RETURN;END_IFaBuffer[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_INPUTsTopicFilter : STRING; // 需要检查的 Topic FilterEND_VARVARuiLen : UINT; // Topic Filter 长度[byte]uiIndex : UINT; // 字符逐字节扫描索引[byte]END_VAR// === IMPLEMENTATION ===uiLen := TO_UINT(LEN(sTopicFilter));IF uiLen = 0 THENF_MqttContainsWildcard := FALSE;RETURN;END_IFFOR uiIndex := 0 TO uiLen - 1 DOIF (sTopicFilter[uiIndex] = 16#2B) OR (sTopicFilter[uiIndex] = 16#23) THENF_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_INPUTuiStartIndex : UINT; // Remaining Length 在缓冲区中的起始索引,通常为 1[byte]uiBufferLen : UINT; // 当前缓冲区内有效数据长度[byte]END_VARVAR_IN_OUTaBuffer : ARRAY[*] OF BYTE; // MQTT 原始接收缓冲区END_VARVAR_OUTPUTudiValue : UDINT; // 解码后的 Remaining Length 数值[byte]uiBytesUsed : UINT; // Remaining Length 字段实际占用的字节数[byte]xNeedMore : BOOL; // 数据不足时置 TRUE,上层应继续读取 TCP 数据END_VARVARudiMultiplier : UDINT; // MQTT 变长整数倍率,依次为 1、128、16384、2097152byEncoded : BYTE; // 当前读取的编码字节uiIndex : UINT; // 当前读取缓冲区索引[byte]uiLoop : UINT; // 变长整数最多 4 字节的循环计数END_VAR// === IMPLEMENTATION ===udiValue := 0;uiBytesUsed := 0;xNeedMore := FALSE;udiMultiplier := 1;uiIndex := uiStartIndex;IF uiStartIndex >= uiBufferLen THENxNeedMore := TRUE;F_MqttDecodeRemainingLength := FALSE;RETURN;END_IFFOR uiLoop := 1 TO 4 DOIF uiIndex >= uiBufferLen THENxNeedMore := TRUE;F_MqttDecodeRemainingLength := FALSE;RETURN;END_IFbyEncoded := aBuffer[uiIndex];udiValue := udiValue + TO_UDINT(byEncoded AND 16#7F) * udiMultiplier;uiBytesUsed := uiBytesUsed + 1;IF (byEncoded AND 16#80) = 0 THENF_MqttDecodeRemainingLength := TRUE;RETURN;END_IFudiMultiplier := 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_INPUTudiValue : UDINT; // 需要编码的 MQTT Remaining Length 数值[byte]udiBufferSize : UDINT; // 调用方输出缓冲区容量[byte]END_VARVAR_IN_OUTaBuffer : ARRAY[*] OF BYTE; // 调用方输出缓冲区,函数从索引 0 起写入编码结果END_VARVAR_OUTPUTuiEncodedLen : UINT; // 实际编码产生的字节数[byte]END_VARVARudiWorkValue : UDINT; // 编码过程中逐步除以 128 的临时值byEncoded : BYTE; // 当前轮生成的 7 位数据和 continuation 标志uiIndex : UINT; // 当前写入缓冲区的索引[byte]END_VAR// === IMPLEMENTATION ===uiEncodedLen := 0;udiWorkValue := udiValue;IF udiValue > GVL_MqttBroker.cnMqttRemainingLengthMax THENF_MqttEncodeRemainingLength := FALSE;RETURN;END_IFREPEATIF TO_UDINT(uiIndex) >= udiBufferSize THENF_MqttEncodeRemainingLength := FALSE;RETURN;END_IFbyEncoded := TO_BYTE(udiWorkValue MOD 128);udiWorkValue := udiWorkValue / 128;IF udiWorkValue > 0 THENbyEncoded := byEncoded OR 16#80;END_IFaBuffer[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_INPUTsFilter : STRING; // 待校验的 MQTT Topic FilterEND_VARVARuiLen : 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 THENF_MqttIsValidTopicFilter := FALSE;RETURN;END_IFIF uiLen > GVL_MqttBroker.cnMaxTopicLen THENF_MqttIsValidTopicFilter := FALSE;RETURN;END_IFFOR uiIndex := 0 TO uiLen - 1 DOIF uiIndex > 0 THENbyPrev := sFilter[uiIndex - 1];ELSEbyPrev := 0;END_IFIF uiIndex < (uiLen - 1) THENbyNext := sFilter[uiIndex + 1];ELSEbyNext := 0;END_IFCASE sFilter[uiIndex] OF16#23:// # 必须是最后一个字符,并且要么单独出现,要么前面是层级分隔符 /。IF uiIndex <> (uiLen - 1) THENF_MqttIsValidTopicFilter := FALSE;RETURN;END_IFIF (uiIndex > 0) AND (byPrev <> 16#2F) THENF_MqttIsValidTopicFilter := FALSE;RETURN;END_IF16#2B:// + 必须独占一个层级,左右只能是边界或层级分隔符 /。IF (uiIndex > 0) AND (byPrev <> 16#2F) THENF_MqttIsValidTopicFilter := FALSE;RETURN;END_IFIF (uiIndex < (uiLen - 1)) AND (byNext <> 16#2F) THENF_MqttIsValidTopicFilter := FALSE;RETURN;END_IFELSE// 普通字符不需要额外限制;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_INPUTsTopic : STRING; // 待校验的 MQTT Topic NameEND_VARVARuiLen : UINT; // Topic Name 长度[byte]uiIndex : UINT; // 当前检查的字符位置,CODESYS STRING 字符下标从 0 开始[byte]END_VAR// === IMPLEMENTATION ===uiLen := TO_UINT(LEN(sTopic));IF uiLen = 0 THENF_MqttIsValidTopicName := FALSE;RETURN;END_IFIF uiLen > GVL_MqttBroker.cnMaxTopicLen THENF_MqttIsValidTopicName := FALSE;RETURN;END_IFFOR uiIndex := 0 TO uiLen - 1 DOIF (sTopic[uiIndex] = 16#2B) OR (sTopic[uiIndex] = 16#23) THENF_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_INPUTuiBufferLen : UINT; // 当前 MQTT 报文有效长度[byte]uiMaxLen : UINT; // 输出字符串允许保存的最大长度[byte]END_VARVAR_IN_OUTaBuffer : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区uiOffset : UINT; // 当前读取偏移,成功后推进到字符串末尾后一字节[byte]sValue : STRING; // 读取出的字符串内容,短字段由调用方使用默认长度临时缓冲转接END_VARVAR_OUTPUTuiStringLen : UINT; // MQTT 字符串字段声明的原始长度[byte]END_VARVARuiIndex : UINT; // 字符串字符复制索引,CODESYS STRING 字符下标从 0 开始[byte]uiReadIndex : UINT; // 当前读取报文缓冲区索引[byte]END_VAR// === IMPLEMENTATION ===sValue := '';uiStringLen := 0;IF (uiOffset + 1) >= uiBufferLen THENF_MqttReadString := FALSE;RETURN;END_IFuiStringLen := TO_UINT(aBuffer[uiOffset]) * 256 + TO_UINT(aBuffer[uiOffset + 1]);uiOffset := uiOffset + 2;IF uiStringLen > uiMaxLen THENF_MqttReadString := FALSE;RETURN;END_IFIF (TO_UDINT(uiOffset) + TO_UDINT(uiStringLen)) > TO_UDINT(uiBufferLen) THENF_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 DOuiReadIndex := uiOffset + uiIndex;IF uiReadIndex >= uiBufferLen THENF_MqttReadString := FALSE;RETURN;END_IFsValue[uiIndex] := aBuffer[uiReadIndex];END_FORsValue[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_INPUTuiBufferLen : UINT; // 当前 MQTT 完整报文长度或可用缓冲长度[byte]END_VARVAR_IN_OUTaBuffer : ARRAY[*] OF BYTE; // MQTT 原始报文缓冲区uiOffset : UINT; // 输入为变长整数起始偏移,成功后推进到变长整数之后[byte]END_VARVAR_OUTPUTudiValue : UDINT; // 解码出的变长整数数值,MQTT 5.0 属性长度使用该值[byte]END_VARVARudiMultiplier : UDINT; // MQTT 变长整数倍率,依次为 1、128、16384、2097152byEncoded : BYTE; // 当前读取的编码字节uiLoop : UINT; // 变长整数最多允许 4 个字节END_VAR// === IMPLEMENTATION ===udiValue := 0;udiMultiplier := 1;IF uiOffset >= uiBufferLen THENF_MqttSkipVariableByteInteger := FALSE;RETURN;END_IFFOR uiLoop := 1 TO 4 DOIF uiOffset >= uiBufferLen THENF_MqttSkipVariableByteInteger := FALSE;RETURN;END_IFbyEncoded := aBuffer[uiOffset];udiValue := udiValue + TO_UDINT(byEncoded AND 16#7F) * udiMultiplier;uiOffset := uiOffset + 1;IF (byEncoded AND 16#80) = 0 THENF_MqttSkipVariableByteInteger := TRUE;RETURN;END_IFudiMultiplier := 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_INPUTsValue : STRING; // 需要检查的完整字符串sPrefix : STRING; // 期望匹配的前缀,空前缀表示全部匹配END_VARVARuiValueLen : 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 THENF_MqttStartsWith := TRUE;RETURN;END_IFIF uiValueLen < uiPrefixLen THENF_MqttStartsWith := FALSE;RETURN;END_IFFOR uiIndex := 0 TO uiPrefixLen - 1 DOIF sValue[uiIndex] <> sPrefix[uiIndex] THENF_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 事务调度、重试和生产级闭环
夜雨聆风