ControlRookie
返回文章

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

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

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

适合谁收藏

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

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

先给结论

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

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

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

从工程问题到代码职责

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

本篇源码覆盖表

序号源码对象行数
1FB_MqttBrokerCodec.M_BuildPublish.st146
2FB_MqttBrokerCodec.M_BuildSimpleAck.st277
3FB_MqttBrokerCodec.M_ParseConnect.st290
4FB_MqttBrokerCodec.M_ParsePublish.st143
5FB_MqttBrokerCodec.M_ParseSubscribe.st145
6FB_MqttBrokerCodec.M_ParseUnsubscribe.st122
7FB_MqttBrokerCodec.st14
8F_MqttAppendString.st49
9F_MqttContainsWildcard.st34
10F_MqttDecodeRemainingLength.st63
11F_MqttEncodeRemainingLength.st55
12F_MqttIsValidTopicFilter.st76
13F_MqttIsValidTopicName.st39
14F_MqttReadString.st67
15F_MqttSkipVariableByteInteger.st54
16F_MqttStartsWith.st42

推荐阅读顺序

  • 先看 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 事务调度、重试和生产级闭环
评论和回复区

评论区预留

这里先保留评论和回复结构,不接入第三方服务。后续统一决定登录、匿名、审核、反垃圾和静态站兼容策略。

↑ ↓