源码加更05_PUBLISH、SUBSCRIBE、UNSUBSCRIBE 业务报文实现
[!abstract] 这一篇完整公开发布、订阅、取消订阅的报文构建和业务响应处理。这是 MQTT Client 从“能连上”走到“能干活”的关键层。
适合谁收藏
想对照 PUBLISH/SUBSCRIBE 报文字段和 ST 代码的读者。 需要排查主题过滤器、订阅表和业务命令边界的工程师。 想知道 QoS0/1/2 在发布报文中如何分流的人。
本篇核心图

读图重点:先看源码对象之间的职责边界,再看数据、状态和错误如何沿着调用链流动。源码加更不是把文件名列出来,而是把完整代码、工程意图和验证口径一起讲清楚。
先给结论
PUBLISH、SUBSCRIBE、UNSUBSCRIBE 都是业务动作,但它们对状态机的要求不一样。发布要管 QoS 和 PacketId,订阅要管 Topic Filter 和 SUBACK,取消订阅要管本地订阅表收口。
从理论到代码实现链路
MQTT 标准给的是报文类型、固定头、可变头、载荷、QoS 交互和会话语义;PLC 工程真正要解决的是周期扫描、缓冲区长度、错误锁存、在线变量、连接重入和现场可诊断性。
所以这套开源实现不能只按协议章节拆,也不能只按文件名拆。正确读法是把标准约束翻译成程序对象:入口程序负责给命令和观测点,GVL 和 DUT 定义容量与数据模型,主功能块负责调度状态机,构建方法负责出站报文,处理方法负责入站报文,辅助方法负责长度、队列、事务、主题和诊断边界。
本篇完整公开业务报文构建、PUBLISH/SUBACK/UNSUBACK 处理和订阅表辅助方法。
再往下一层看,这里其实有两条线同时存在。第一条是协议线:固定头、Remaining Length、PacketId、QoS、Topic、Payload 和 Reason Code 必须能按 MQTT 规则组合起来。第二条是 PLC 工程线:每个周期只能推进有限步骤,所有中间状态都要能被在线变量观察,所有错误都要能被锁存并归类,所有缓冲区长度都要在写入前被检查。
这就是源码加更必须完整公开的原因。只给几段核心片段,读者最多能看懂某个判断;把完整对象放出来,读者才能看到对象之间如何传递状态、长度、错误和诊断信息。完整源码讲解不是为了堆代码,而是为了让读者能从标准约束一路追到可运行的 ST 对象,再从现场现象反向定位到具体边界。
本篇公开的完整源码范围
M_BuildPublishPacket.st | ||
M_BuildSubscribePacket.st | ||
M_BuildUnsubscribePacket.st | ||
M_HandlePublish.st | ||
M_HandleSubAck.st | ||
M_HandleUnsubAck.st | ||
M_IsValidTopicFilter.st | ||
M_SubListAdd.st | ||
M_SubListClear.st | ||
M_SubListRemove.st |
怎么读这些源码
第一遍只看对象职责:这个文件解决哪一层问题,是入口、模型、状态、构建、接收、事务,还是诊断。
第二遍看边界变量:长度、索引、PacketId、QoS、状态枚举、错误码、缓冲区水位和在线观测量。PLC 通信代码最怕的是“能跑但不可诊断”,所以每个关键对象都要问一句:现场出问题时,我能不能从它留下的变量看出原因。
第三遍再看具体语句。源码全部公开,不等于读者要从第一行顺序读到最后一行。更稳的方式是用图和表先建立地图,再回到完整代码里确认每个边界确实落地。
工程验证路径
验证时看三组量:出站报文长度是否正确,SUBACK/UNSUBACK 是否更新本地订阅表,PUBLISH 接收后是否锁存主题和载荷。
本篇完整开源代码
完整代码 1:M_BuildPublishPacket.st
这一段完整公开 M_BuildPublishPacket.st。读代码时先看对象职责,再看状态、长度、错误和返回值,不要只抄几行赋值。
/// =======================================================================/// 名称 : M_BuildPublishPacket/// 功能 : 构建 PUBLISH 发送报文/// 说明 : 根据主题、载荷、QoS 和协议版本组装发布报文。/// 编程人员 : ControlRookie/// 时间 : 2026-05-05/// 版本 : V1.1/// ======================================================================={attribute 'hide_all_locals'}METHOD M_BuildPublishPacket : BOOLVARuiPos : UINT := 0; // 当前写入发送缓冲区的位置偏移[byte]uiVarHeaderLen : UINT; // PUBLISH 可变报头总长度[byte]uiPayloadLen : UINT; // 当前待发布载荷长度[byte]uiRemainingLen : UINT; // 写入固定报头中的 Remaining Length 值[byte]uiPropsLen : UINT; // MQTT 5.0 PUBLISH 属性区总长度(含属性长度字段)[byte]uiPropertyDataLen : UINT; // MQTT 5.0 PUBLISH 属性内容长度(不含属性长度字段)[byte]uiTopicAlias : UINT; // 本次准备写入报文的主题别名编号uiInflightIndex : UINT; // 在途队列中登记或重发的槽位索引uiPublishPacketId : UINT; // 本次发布实际使用的 Packet Identifieri : DINT; // 清空缓冲区或扫描主题时使用的循环索引sPublishTopic : STRING(GVL_Mqtt.cnMaxTopicLen); // 本次真正写入报文的主题字符串sPublishPayload : STRING(GVL_Mqtt.cnMaxPayloadSize);// 本次真正写入报文的载荷字符串ePublishQoS : E_MqttQoS; // 本次真正写入报文的 QoS 等级bPublishRetainLocal : BOOL; // 本次真正写入报文的 Retain 标志bHasWildcard : BOOL; // 发布主题中是否误带通配符END_VAR// === IMPLEMENTATION ===/// 先做最基础的资源与配额保护:/// - 发送缓冲区太小就不允许继续构包;/// - MQTT 5.0 下 Send Quota 已耗尽时,禁止再发新的 QoS>0 报文。IF SIZEOF(aTxBuf) < 256 THENM_BuildPublishPacket := FALSE;RETURN;END_IFIF (ePubQoS > E_MqttQoS.byQoS0) AND (eVersion = E_MqttVersion.byMqttVersion50) AND (uiSendQuota = 0) THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrReceiveMaxExceeded),sMessage := 'Send quota exhausted');M_BuildPublishPacket := FALSE;RETURN;END_IFFOR i := LOWER_BOUND(aTxBuf, 1) TO UPPER_BOUND(aTxBuf, 1) DOaTxBuf[i] := 0;END_FOR/// 默认使用当前输入引脚的发布参数;/// 如果这次是 inflight 超时重发,则整包参数改为从在途槽位中恢复。sPublishTopic := sPubTopic;sPublishPayload := sPubPayload;ePublishQoS := ePubQoS;bPublishRetainLocal := bPubRetain;uiPublishPacketId := 0;IF (uiRetryInflightIndex > 0) AND (uiRetryInflightIndex <= GVL_Mqtt.cnMaxInflight) THENIF aInflight[uiRetryInflightIndex].bUsed THENsPublishTopic := aInflight[uiRetryInflightIndex].sTopic;sPublishPayload := aInflight[uiRetryInflightIndex].sPayload;ePublishQoS := aInflight[uiRetryInflightIndex].eQoS;bPublishRetainLocal := aInflight[uiRetryInflightIndex].bRetain;uiPublishPacketId := aInflight[uiRetryInflightIndex].uiPacketId;bDup := aInflight[uiRetryInflightIndex].bDup;END_IFEND_IF/// 先做主题、载荷、UTF-8、通配符等基础合法性校验,/// 这些都属于“构包前就必须被本地拒绝”的错误。IF sPublishTopic = '' THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),sMessage := 'Publish topic is required');M_BuildPublishPacket := FALSE;RETURN;END_IFIF TO_UINT(LEN(sPublishTopic)) > GVL_Mqtt.cnMaxTopicLen THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),sMessage := 'Publish topic exceeds maximum length');M_BuildPublishPacket := FALSE;RETURN;END_IFIF NOT M_IsValidUtf8String(sValue := sPublishTopic) THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),sMessage := 'Publish topic is not valid UTF-8');M_BuildPublishPacket := FALSE;RETURN;END_IFbHasWildcard := FALSE;FOR i := 1 TO TO_DINT(LEN(sPublishTopic)) DOIF (sPublishTopic[i - 1] = 16#2B) OR (sPublishTopic[i - 1] = 16#23) THENbHasWildcard := TRUE;EXIT;END_IFEND_FORIF bHasWildcard THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),sMessage := 'Publish topic must not contain wildcards');M_BuildPublishPacket := FALSE;RETURN;END_IFIF TO_UINT(LEN(sPublishPayload)) > GVL_Mqtt.cnMaxPayloadSize THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrPacketTooLarge),sMessage := 'Publish payload exceeds configured maximum');M_BuildPublishPacket := FALSE;RETURN;END_IFIF eVersion = E_MqttVersion.byMqttVersion50 THEN/// MQTT 5.0 下,本地还要受服务端能力约束:/// 例如最大 QoS、是否支持 Retain、是否限制最大报文长度。IF TO_BYTE(ePublishQoS) > byServerMaxQoS THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),sMessage := 'Publish QoS exceeds server limit');M_BuildPublishPacket := FALSE;RETURN;END_IFIF bPublishRetainLocal AND (NOT bServerRetainAvailable) THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),sMessage := 'Retain not supported by server');M_BuildPublishPacket := FALSE;RETURN;END_IFEND_IFIF ePublishQoS = E_MqttQoS.byQoS0 THENuiVarHeaderLen := 2 + TO_UINT(LEN(sPublishTopic));ELSEuiVarHeaderLen := 2 + TO_UINT(LEN(sPublishTopic)) + 2;END_IFuiPayloadLen := TO_UINT(LEN(sPublishPayload));uiRemainingLen := uiVarHeaderLen + uiPayloadLen;IF eVersion = E_MqttVersion.byMqttVersion50 THEN/// 当前 V2.0 的 PUBLISH 只在需要时附加 Topic Alias 属性,/// 其余属性未来可继续在这里扩展,而不会污染 3.1.1 主链路。uiPropertyDataLen := 0;uiTopicAlias := 0;IF uiServerTopicAliasMax > 0 THENuiTopicAlias := uiNextTopicAlias;IF uiTopicAlias = 0 THENuiTopicAlias := 1;END_IFuiPropertyDataLen := uiPropertyDataLen + 3;END_IFuiPropsLen := 1 + uiPropertyDataLen;uiRemainingLen := uiRemainingLen + uiPropsLen;IF (udServerMaxPacketSize > 0) AND (TO_UDINT(uiRemainingLen + 5) > udServerMaxPacketSize) THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrPacketTooLarge),sMessage := 'Publish packet exceeds server maximum packet size');M_BuildPublishPacket := FALSE;RETURN;END_IFEND_IFIF uiRemainingLen + 5 > SIZEOF(aTxBuf) THENM_BuildPublishPacket := FALSE;RETURN;END_IF/// QoS0 不允许带 DUP;QoS1/2 则按当前重发状态决定是否置位 DUP。IF ePublishQoS = E_MqttQoS.byQoS0 THENbDup := FALSE;END_IF/// 先写固定报头,再写 Remaining Length、主题字符串、Packet ID、属性区和 payload。aTxBuf[uiPos] := E_MqttPacketType.byPublish OR(BOOL_TO_BYTE(bDup) * 16#08) ORTO_BYTE(SHL(ePublishQoS AND 16#03, 1)) OR(BOOL_TO_BYTE(bPublishRetainLocal) * 16#01);uiPos := uiPos + 1;uiPos := uiPos + M_EncodeRemainingLength(udiLength := uiRemainingLen, pBuffer := ADR(aTxBuf[uiPos]));uiPos := uiPos + M_AppendString(sStr := sPublishTopic, pBuffer := ADR(aTxBuf[uiPos]));IF ePublishQoS > E_MqttQoS.byQoS0 THEN/// 新发 QoS1/2 报文需要申请新的 Packet Identifier;/// 重发报文则沿用 inflight 里原来的 Packet Identifier。IF uiPublishPacketId = 0 THENuiExpectedPacketId := M_GetNextPacketId();ELSEuiExpectedPacketId := uiPublishPacketId;END_IFIF uiExpectedPacketId = 0 THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrReceiveMaxExceeded),sMessage := 'No Packet ID available');M_BuildPublishPacket := FALSE;RETURN;END_IFuiQoS2PacketId := uiExpectedPacketId;aTxBuf[uiPos] := UINT_TO_BYTE(SHR(uiExpectedPacketId, 8));uiPos := uiPos + 1;aTxBuf[uiPos] := UINT_TO_BYTE(uiExpectedPacketId AND 16#FF);uiPos := uiPos + 1;END_IFIF eVersion = E_MqttVersion.byMqttVersion50 THEN/// Topic Alias 采用“循环递增直到服务端允许上限”的轻量策略。uiPos := uiPos + M_EncodeRemainingLength(udiLength := uiPropertyDataLen, pBuffer := ADR(aTxBuf[uiPos]));IF uiTopicAlias > 0 THENaTxBuf[uiPos] := GVL_Mqtt.cnPropTopicAlias;uiPos := uiPos + 1;aTxBuf[uiPos] := UINT_TO_BYTE(SHR(uiTopicAlias, 8));uiPos := uiPos + 1;aTxBuf[uiPos] := UINT_TO_BYTE(uiTopicAlias AND 16#FF);uiPos := uiPos + 1;uiNextTopicAlias := uiTopicAlias + 1;IF (uiNextTopicAlias = 0) OR (uiNextTopicAlias > uiServerTopicAliasMax) THENuiNextTopicAlias := 1;END_IFEND_IFEND_IFuiPos := uiPos + M_AppendPayload(sPayload := sPublishPayload, pBuffer := ADR(aTxBuf[uiPos]));uiTxLength := uiPos;IF ePublishQoS > E_MqttQoS.byQoS0 THEN/// 只有新的 QoS1/2 发布才需要入 inflight 队列;/// 重发场景只更新时间戳,不重复创建槽位。IF uiPublishPacketId = 0 THENuiInflightIndex := M_InflightAdd(uiPacketId := uiExpectedPacketId,eQoS := ePublishQoS,sTopic := sPublishTopic,sPayload := sPublishPayload,uiPayloadLen := uiPayloadLen,bRetain := bPublishRetainLocal);IF uiInflightIndex = 0 THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrReceiveMaxExceeded),sMessage := 'Inflight queue is full');M_BuildPublishPacket := FALSE;RETURN;END_IFELSIF (uiRetryInflightIndex > 0) AND (uiRetryInflightIndex <= GVL_Mqtt.cnMaxInflight) THENaInflight[uiRetryInflightIndex].tLastSend := TIME();END_IFEND_IFM_BuildPublishPacket := TRUE;
完整代码 2:M_BuildSubscribePacket.st
这一段完整公开 M_BuildSubscribePacket.st。读代码时先看对象职责,再看状态、长度、错误和返回值,不要只抄几行赋值。
/// =======================================================================/// 名称 : M_BuildSubscribePacket/// 功能 : 构建 SUBSCRIBE 发送报文/// 说明 : 根据 MQTT 版本组装订阅报文,并在发送前执行基础协议校验。/// 编程人员 : ControlRookie/// 时间 : 2026-05-05/// 版本 : V1.2/// ======================================================================={attribute 'hide_all_locals'}METHOD M_BuildSubscribePacket : BOOLVARuiPos : UINT := 0; // 当前写入发送缓冲区的位置偏移[byte]uiVarHeaderLen : UINT; // SUBSCRIBE 可变报头总长度[byte]uiPayloadLen : UINT; // 本次订阅主题过滤器载荷总长度[byte]uiRemainingLen : UINT; // 写入固定报头中的 Remaining Length 值[byte]uiPropsLen : UINT; // MQTT 5.0 SUBSCRIBE 属性区总长度[byte]uiSubIdVbiBytes : UINT; // 订阅标识符编码成 VBI 后实际占用的字节数[byte]i : DINT; // 扫描主题字符串或清空缓冲区时使用的循环索引bHasWildcard : BOOL; // 当前订阅主题过滤器是否使用了 + / # 通配符bIsSharedSubscription : BOOL; // 当前订阅主题过滤器是否采用 share/` 前缀代表共享订阅。/// 这里只做最小前缀识别,真正主题过滤器合法性已在前面统一校验。IF (sActiveSubTopic[0] = 16#24) AND(sActiveSubTopic[1] = 16#73) AND(sActiveSubTopic[2] = 16#68) AND(sActiveSubTopic[3] = 16#61) AND(sActiveSubTopic[4] = 16#72) AND(sActiveSubTopic[5] = 16#65) AND(sActiveSubTopic[6] = 16#2F) THENbIsSharedSubscription := TRUE;END_IFEND_IFIF (eVersion = E_MqttVersion.byMqttVersion50) AND (udiActiveSubscriptionId > 0) AND (NOT bServerSubIdAvail) THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),sMessage := 'Server does not support subscription identifiers');M_BuildSubscribePacket := FALSE;RETURN;END_IFIF (eVersion = E_MqttVersion.byMqttVersion50) AND (NOT bServerWildcardSubAvail) AND bHasWildcard THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),sMessage := 'Server does not support wildcard subscriptions');M_BuildSubscribePacket := FALSE;RETURN;END_IFIF (eVersion = E_MqttVersion.byMqttVersion50) AND (NOT bServerSharedSubAvail) AND bIsSharedSubscription THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),sMessage := 'Server does not support shared subscriptions');M_BuildSubscribePacket := FALSE;RETURN;END_IFFOR i := LOWER_BOUND(aTxBuf, 1) TO UPPER_BOUND(aTxBuf, 1) DOaTxBuf[i] := 0;END_FOR/// SUBSCRIBE 可变报头至少包含 Packet Identifier;/// MQTT 5.0 下如果启用了 Subscription Identifier,还要额外预留属性长度。uiVarHeaderLen := 2;uiPropsLen := 0;uiSubIdVbiBytes := 0;IF eVersion = E_MqttVersion.byMqttVersion50 THENIF udiActiveSubscriptionId > 0 THENIF udiActiveSubscriptionId < 128 THENuiSubIdVbiBytes := 1;ELSIF udiActiveSubscriptionId < 16384 THENuiSubIdVbiBytes := 2;ELSIF udiActiveSubscriptionId < 2097152 THENuiSubIdVbiBytes := 3;ELSEuiSubIdVbiBytes := 4;END_IFuiPropsLen := uiPropsLen + 1 + uiSubIdVbiBytes;END_IFIF uiPropsLen < 128 THENuiVarHeaderLen := uiVarHeaderLen + 1 + uiPropsLen;ELSEuiVarHeaderLen := uiVarHeaderLen + 2 + uiPropsLen;END_IFEND_IFuiPayloadLen := TO_UINT(LEN(sActiveSubTopic)) + 3;uiRemainingLen := uiVarHeaderLen + uiPayloadLen;IF uiRemainingLen + 5 > SIZEOF(aTxBuf) THENM_BuildSubscribePacket := FALSE;RETURN;END_IFuiPos := 0;aTxBuf[uiPos] := E_MqttPacketType.bySubscribe;uiPos := uiPos + 1;uiPos := uiPos + M_EncodeRemainingLength(udiLength := uiRemainingLen, pBuffer := ADR(aTxBuf[uiPos]));/// 每一笔订阅事务都必须独占一个 Packet Identifier,/// 后续 SUBACK 就靠它来和“当前激活订阅请求”做精确匹配。uiPendingSubPacketId := M_GetNextPacketId();IF uiPendingSubPacketId = 0 THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrReceiveMaxExceeded),sMessage := 'No Packet ID available for subscribe');M_BuildSubscribePacket := FALSE;RETURN;END_IFuiExpectedPacketId := uiPendingSubPacketId;aTxBuf[uiPos] := UINT_TO_BYTE(SHR(uiPendingSubPacketId, 8));uiPos := uiPos + 1;aTxBuf[uiPos] := UINT_TO_BYTE(uiPendingSubPacketId AND 16#FF);uiPos := uiPos + 1;IF eVersion = E_MqttVersion.byMqttVersion50 THEN/// 当前 V2.0 的 SUBSCRIBE 属性区只编码 Subscription Identifier;/// 其他订阅选项留作后续 5.0 扩展项。uiPos := uiPos + M_EncodeRemainingLength(udiLength := uiPropsLen, pBuffer := ADR(aTxBuf[uiPos]));IF udiActiveSubscriptionId > 0 THENaTxBuf[uiPos] := GVL_Mqtt.cnPropSubscriptionId;uiPos := uiPos + 1;uiPos := uiPos + M_EncodeRemainingLength(udiLength := udiActiveSubscriptionId, pBuffer := ADR(aTxBuf[uiPos]));END_IFEND_IFuiPos := uiPos + M_AppendString(sStr := sActiveSubTopic, pBuffer := ADR(aTxBuf[uiPos]));/// 当前订阅选项字节只写入 QoS;/// No Local / Retain As Published / Retain Handling 还没有进入 V2.0 当前主线。aTxBuf[uiPos] := TO_BYTE(eActiveSubQoS AND 16#03);uiPos := uiPos + 1;uiTxLength := uiPos;M_BuildSubscribePacket := TRUE;
完整代码 3:M_BuildUnsubscribePacket.st
这一段完整公开 M_BuildUnsubscribePacket.st。读代码时先看对象职责,再看状态、长度、错误和返回值,不要只抄几行赋值。
/// =======================================================================/// 名称 : M_BuildUnsubscribePacket/// 功能 : 构建 UNSUBSCRIBE 发送报文/// 说明 : 根据 MQTT 版本组装取消订阅报文并更新发送长度。/// 编程人员 : ControlRookie/// 时间 : 2026-05-05/// 版本 : V1.0/// ======================================================================={attribute 'hide_all_locals'}METHOD M_BuildUnsubscribePacket : BOOLVARuiPos : UINT := 0; // 当前写入发送缓冲区的位置偏移[byte]uiVarHeaderLen : UINT; // UNSUBSCRIBE 可变报头总长度[byte]uiPayloadLen : UINT; // 取消订阅主题载荷总长度[byte]uiRemainingLen : UINT; // 写入固定报头中的 Remaining Length 值[byte]uiPropsLen : UINT; // MQTT 5.0 UNSUBSCRIBE 属性区总长度[byte]i : DINT; // 清空发送缓冲区时使用的循环索引END_VAR// === IMPLEMENTATION ===// BUG-10: 缓冲区溢出保护IF SIZEOF(aTxBuf) < 256 THENM_BuildUnsubscribePacket := FALSE;RETURN;END_IFIF sUnsubTopic = '' THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),sMessage := 'Unsubscribe topic is required');M_BuildUnsubscribePacket := FALSE;RETURN;END_IFIF TO_UINT(LEN(sUnsubTopic)) > GVL_Mqtt.cnMaxTopicLen THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),sMessage := 'Unsubscribe topic exceeds maximum length');M_BuildUnsubscribePacket := FALSE;RETURN;END_IFIF NOT M_IsValidUtf8String(sValue := sUnsubTopic) THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),sMessage := 'Unsubscribe topic filter is not valid UTF-8');M_BuildUnsubscribePacket := FALSE;RETURN;END_IFIF NOT M_IsValidTopicFilter(sTopicFilter := sUnsubTopic) THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),sMessage := 'Unsubscribe topic filter is invalid');M_BuildUnsubscribePacket := FALSE;RETURN;END_IF// 清空发送缓冲区,避免数据干扰FOR i := LOWER_BOUND(aTxBuf, 1) TO UPPER_BOUND(aTxBuf, 1) DOaTxBuf[i] := 0;END_FOR/// =======================================================================/// 长度计算/// =======================================================================// 可变报头长度 = Packet ID(2)uiVarHeaderLen := 2;// MQTT 5.0: 取消订阅属性(当前无属性,长度为0)uiPropsLen := 0;IF eVersion = E_MqttVersion.byMqttVersion50 THEN// 属性长度(0)的VBI编码 = 1字节uiVarHeaderLen := uiVarHeaderLen + 1 + uiPropsLen;END_IF// 取消订阅主题载荷长度 = 主题长度前缀 2 字节 + Topic Filter 内容uiPayloadLen := TO_UINT(LEN(sUnsubTopic) + 2);uiRemainingLen := uiVarHeaderLen + uiPayloadLen;// BUG-10: 缓冲区长度检查IF uiRemainingLen + 5 > SIZEOF(aTxBuf) THENM_BuildUnsubscribePacket := FALSE;RETURN;END_IF/// =======================================================================/// 创建报文/// =======================================================================uiPos := 0;// ****************** 固定报文头 = 报文类型 + Remaining Length ******************aTxBuf[0] := E_MqttPacketType.byUnsubscribe;uiPos := uiPos + 1;uiPos := uiPos + M_EncodeRemainingLength(uiRemainingLen, ADR(aTxBuf[uiPos]));// ****************** 可变报文头: Packet ID ******************uiPendingUnsubPacketId := M_GetNextPacketId();IF uiPendingUnsubPacketId = 0 THENM_SetError(uiErrorCode := TO_UINT(E_ReasonCode.uiErrReceiveMaxExceeded),sMessage := 'No Packet ID available for unsubscribe');M_BuildUnsubscribePacket := FALSE;RETURN;END_IFuiExpectedPacketId := uiPendingUnsubPacketId;aTxBuf[uiPos] := UINT_TO_BYTE(SHR(uiPendingUnsubPacketId, 8)); uiPos := uiPos + 1;aTxBuf[uiPos] := UINT_TO_BYTE(uiPendingUnsubPacketId AND 16#FF); uiPos := uiPos + 1;// ****************** MQTT 5.0: 取消订阅属性 ******************IF eVersion = E_MqttVersion.byMqttVersion50 THEN// 属性长度(VBI编码,当前为0=无属性)uiPos := uiPos + M_EncodeRemainingLength(uiPropsLen, ADR(aTxBuf[uiPos]));END_IF// ****************** 载荷: 主题 ******************uiPos := uiPos + M_AppendString(sUnsubTopic, ADR(aTxBuf[uiPos]));uiTxLength := uiPos;M_BuildUnsubscribePacket := TRUE;
完整代码 4:M_HandlePublish.st
这一段完整公开 M_HandlePublish.st。读代码时先看对象职责,再看状态、长度、错误和返回值,不要只抄几行赋值。
/// =======================================================================/// 名称 : M_HandlePublish/// 功能 : 兼容桩方法,保留 PUBLISH 旧接口/// 说明 : 该方法已废弃,实际处理逻辑已迁移至 M_ProcessReceive,仅为兼容旧调用保留。/// 编程人员 : ControlRookie/// 时间 : 2026-05-05/// 版本 : V1.0/// ======================================================================={attribute 'hide_all_locals'}METHOD M_HandlePublish : BOOL// === IMPLEMENTATION ===// 逻辑已移至M_ProcessReceiveM_HandlePublish := FALSE;
完整代码 5:M_HandleSubAck.st
这一段完整公开 M_HandleSubAck.st。读代码时先看对象职责,再看状态、长度、错误和返回值,不要只抄几行赋值。
/// =======================================================================/// 名称 : M_HandleSubAck/// 功能 : 兼容桩方法,保留 SUBACK 旧接口/// 说明 : 该方法已废弃,实际处理逻辑已迁移至 M_ProcessReceive,仅为兼容旧调用保留。/// 编程人员 : ControlRookie/// 时间 : 2026-05-05/// 版本 : V1.0/// ======================================================================={attribute 'hide_all_locals'}METHOD M_HandleSubAck : BOOLVAR_INPUTEND_VAR// === IMPLEMENTATION ===// 逻辑已移至M_ProcessReceiveM_HandleSubAck := FALSE;
完整代码 6:M_HandleUnsubAck.st
这一段完整公开 M_HandleUnsubAck.st。读代码时先看对象职责,再看状态、长度、错误和返回值,不要只抄几行赋值。
/// =======================================================================/// 名称 : M_HandleUnsubAck/// 功能 : 兼容桩方法,保留 UNSUBACK 旧接口/// 说明 : 该方法已废弃,实际处理逻辑已迁移至 M_ProcessReceive,仅为兼容旧调用保留。/// 编程人员 : ControlRookie/// 时间 : 2026-05-05/// 版本 : V1.0/// ======================================================================={attribute 'hide_all_locals'}METHOD M_HandleUnsubAck : BOOLVAR_INPUTEND_VAR// === IMPLEMENTATION ===// 逻辑已移至M_ProcessReceiveM_HandleUnsubAck := FALSE;
完整代码 7:M_IsValidTopicFilter.st
这一段完整公开 M_IsValidTopicFilter.st。读代码时先看对象职责,再看状态、长度、错误和返回值,不要只抄几行赋值。
/// =======================================================================/// 名称 : M_IsValidTopicFilter/// 功能 : 校验 MQTT Topic Filter 合法性/// 说明 : 用于 SUBSCRIBE 主题过滤器的离线格式校验/// 编程人员 : ControlRookie/// 时间 : 2026-05-05/// 版本 : V1.0/// ======================================================================={attribute 'hide_all_locals'}METHOD M_IsValidTopicFilter : BOOLVAR_INPUTsTopicFilter : STRING(GVL_Mqtt.cnMaxTopicLen); // 待校验主题过滤器END_VARVARuiLen : UINT; // 过滤器长度[char]i : UINT; // 遍历索引byChar : BYTE; // 当前字符byPrev : BYTE; // 前一字符byNext : BYTE; // 后一字符END_VAR// === IMPLEMENTATION ===// Topic Filter 不能为空;空字符串既不是普通主题,也不是合法通配符。uiLen := TO_UINT(LEN(sTopicFilter));IF uiLen = 0 THENM_IsValidTopicFilter := FALSE;RETURN;END_IFFOR i := 1 TO uiLen DObyChar := TO_BYTE(sTopicFilter[i - 1]);// 过滤器中间不允许出现提前结束的 0 字节。IF byChar = 0 THENM_IsValidTopicFilter := FALSE;RETURN;END_IF// '#' 只能出现在最后一层,并且它前面如果还有内容,上一字符必须是 '/'。IF byChar = 16#23 THENIF i <> uiLen THENM_IsValidTopicFilter := FALSE;RETURN;END_IFIF (i > 1) AND (TO_BYTE(sTopicFilter[i - 2]) <> 16#2F) THENM_IsValidTopicFilter := FALSE;RETURN;END_IFEND_IF// '+' 必须独占一层,也就是它左右两边只能是分隔符 '/' 或字符串边界。IF byChar = 16#2B THENbyPrev := 0;byNext := 0;IF i > 1 THENbyPrev := TO_BYTE(sTopicFilter[i - 2]);END_IFIF i < uiLen THENbyNext := TO_BYTE(sTopicFilter[i]);END_IFIF ((i > 1) AND (byPrev <> 16#2F)) OR((i < uiLen) AND (byNext <> 16#2F)) THENM_IsValidTopicFilter := FALSE;RETURN;END_IFEND_IFEND_FOR// 以 '$' 开头的系统主题不能直接写成 '$#' 或 '$+' 这种“跨系统主题根”的订阅方式。IF (uiLen >= 2) AND (sTopicFilter[0] = 16#24) THENIF (sTopicFilter[1] = 16#23) OR (sTopicFilter[1] = 16#2B) THENM_IsValidTopicFilter := FALSE;RETURN;END_IFEND_IFM_IsValidTopicFilter := TRUE;
完整代码 8:M_SubListAdd.st
这一段完整公开 M_SubListAdd.st。读代码时先看对象职责,再看状态、长度、错误和返回值,不要只抄几行赋值。
/// =======================================================================/// 名称 : M_SubListAdd/// 功能 : 将订阅主题加入本地列表/// 说明 : 已存在主题则更新 QoS,不存在时写入空闲槽位。/// 编程人员 : ControlRookie/// 时间 : 2026-05-05/// 版本 : V1.0/// ======================================================================={attribute 'hide_all_locals'}METHOD M_SubListAdd : BOOLVAR_INPUTsTopic : STRING(GVL_Mqtt.cnMaxTopicLen); // 需要登记到本地订阅表的主题过滤器eQos : E_MqttQoS; // 该主题当前已协商成功的订阅 QoS 等级udiSubscriptionId : UDINT; // MQTT 5.0 订阅标识符,后续匹配消息来源时使用END_VARVARi : UINT; // 扫描或写入订阅表时使用的槽位索引xFound : BOOL := FALSE;// 是否已在本地订阅表中找到同名主题END_VAR// === IMPLEMENTATION ===// 检查是否已存在// 同名主题已存在时,不重复新增槽位,只更新本次协商得到的 QoS 和订阅标识符。FOR i := 1 TO GVL_Mqtt.cnMaxSubscriptions DOIF aSubscriptions[i].bActive AND aSubscriptions[i].sTopic = sTopic THENaSubscriptions[i].eQos := eQos;aSubscriptions[i].udiSubscriptionId := udiSubscriptionId;xFound := TRUE;EXIT;END_IFEND_FOR// 添加新订阅// 本地订阅镜像表找空槽位写入,供断线补订和主题统计使用。IF NOT xFound THENFOR i := 1 TO GVL_Mqtt.cnMaxSubscriptions DOIF NOT aSubscriptions[i].bActive THENaSubscriptions[i].sTopic := sTopic;aSubscriptions[i].eQos := eQos;aSubscriptions[i].bActive := TRUE;aSubscriptions[i].udiSubscriptionId := udiSubscriptionId;aSubscriptions[i].bNoLocal := FALSE;aSubscriptions[i].bRetainAsPublished := FALSE;aSubscriptions[i].byRetainHandling := 0;uiSubscriptionCount := uiSubscriptionCount + 1;EXIT;END_IFEND_FOREND_IFM_SubListAdd := TRUE;
完整代码 9:M_SubListClear.st
这一段完整公开 M_SubListClear.st。读代码时先看对象职责,再看状态、长度、错误和返回值,不要只抄几行赋值。
/// =======================================================================/// 名称 : M_SubListClear/// 功能 : 清空订阅主题表/// 说明 : 复位所有订阅槽位与订阅计数。/// 编程人员 : ControlRookie/// 时间 : 2026-05-05/// 版本 : V1.0/// ======================================================================={attribute 'hide_all_locals'}METHOD M_SubListClear : BOOLVARi : UINT; // 逐项清空订阅表时使用的槽位索引END_VAR// === IMPLEMENTATION ===// 断线清理或停机时,本地订阅镜像表需要整体复位,避免残留旧主题快照。FOR i := 1 TO GVL_Mqtt.cnMaxSubscriptions DOaSubscriptions[i].bActive := FALSE;aSubscriptions[i].sTopic := '';aSubscriptions[i].eQoS := E_MqttQoS.byQoS0;aSubscriptions[i].bNoLocal := FALSE;aSubscriptions[i].bRetainAsPublished := FALSE;aSubscriptions[i].byRetainHandling := 0;aSubscriptions[i].udiSubscriptionId := 0;END_FORuiSubscriptionCount := 0;M_SubListClear := TRUE;
完整代码 10:M_SubListRemove.st
这一段完整公开 M_SubListRemove.st。读代码时先看对象职责,再看状态、长度、错误和返回值,不要只抄几行赋值。
/// =======================================================================/// 名称 : M_SubListRemove/// 功能 : 将主题从订阅列表移除/// 说明 : 按主题匹配并清理对应槽位,同时更新订阅计数。/// 编程人员 : ControlRookie/// 时间 : 2026-05-05/// 版本 : V1.0/// ======================================================================={attribute 'hide_all_locals'}METHOD M_SubListRemove : BOOLVAR_INPUTsTopic : STRING(GVL_Mqtt.cnMaxTopicLen); // 需要从本地订阅表中移除的主题过滤器END_VARVARi : UINT; // 扫描订阅表并定位待删除主题时使用的槽位索引END_VAR// === IMPLEMENTATION ===FOR i := 1 TO GVL_Mqtt.cnMaxSubscriptions DOIF aSubscriptions[i].bActive AND aSubscriptions[i].sTopic = sTopic THENaSubscriptions[i].bActive := FALSE;aSubscriptions[i].sTopic := '';aSubscriptions[i].udiSubscriptionId := 0;aSubscriptions[i].bNoLocal := FALSE;aSubscriptions[i].bRetainAsPublished := FALSE;aSubscriptions[i].byRetainHandling := 0;IF uiSubscriptionCount > 0 THENuiSubscriptionCount := uiSubscriptionCount - 1;END_IFEXIT;END_IFEND_FORM_SubListRemove := TRUE;
这一篇你最该记住的几句话
源码加更不是片段展示,而是完整源码对象公开讲解。 先建立对象地图,再读状态、报文和事务,现场调试才不会迷路。 判断源码成熟度,不只看功能是否实现,还要看边界、错误和在线观测量是否闭环。
系列导航
系列定位:MqttClient 系列教程,源码加更阶段,第 15 篇 / 共 16 篇 上一篇:源码加更04 下一篇:源码加更06
夜雨聆风