为 Metatrader 5 开发 MQTT 客户端:TDD 方法 - 第 6 部分·进阶篇
📘

为 Metatrader 5 开发 MQTT 客户端:TDD 方法 - 第 6 部分·进阶篇

第 2/2 篇

「MQTT 发布包的标志位与主题校验逻辑」

在 MT5 里手搓 MQTT 发布报文时,发布标志字节(publish flags)决定了 broker 如何对待这条消息。RETAIN、QoS1、QoS2、DUP 四个标志都通过位运算写进同一个 m_pubflags 变量,例如 SetRetain 里用 retain ? m_pubflags |= RETAIN_FLAG : m_pubflags &= ~RETAIN_FLAG 来置位或清位,SetQoS_1 / SetQoS_2 / SetDup 同理。 QoS 大于 0 时报文必须带 packet ID。代码里的判断很直白:if((m_pubflags & 0x06) != 0) 就调用 SetPacketID,因为 0x06 正好覆盖 QoS1_FLAG 与 QoS2_FLAG 两位。单元测试 TEST_Ctor_Retain_QoS1_TopicName1Char 构造了主题名 "a"、retain=true、QoS1=true 的包,期望字节流是 {51,6,0,1,'a',0,1,0}——首字节 51 即 0x33,是 MQTT PUBLISH 固定头加 RETAIN+QoS1 置位的结果,后面 0,1 就是两字节 packet ID。

主题名不是随便填。SetTopicName 会先拦掉通配符和空字符串:HasWildcardChar(topic_name)StringLen(topic_name) == 0 时直接 ArrayFree 并返回,不编码。正常情况才走 EncodeUTF8String 写进 m_topname。开 MT5 把这段类方法挂到 EA 里,改一个主题名长度或 QoS 组合,用 Build 跑出 pkt 数组对照上面的期望字节,能快速验证你自己的封装有没有偏位。
MQL5 / C++
class="type">void SetCorrelationData(class="type">uchar &binary_data[]);
class="type">void SetUserProperty(const class="type">class="kw">string key, const class="type">class="kw">string val);
class="type">void SetSubscriptionIdentifier(class="type">uint subscript_id);
class="type">void SetContentType(const class="type">class="kw">string content_type);
class=class="str">"cmt">//--- method for setting the payload
class="type">void SetPayload(const class="type">class="kw">string payload);
class=class="str">"cmt">//--- method for building the final packet
class="type">void Build(class="type">uchar &result[]);
};
class="type">bool TEST_Ctor_Retain_QoS1_TopicName1Char()
  {
   Print(__FUNCTION__);
   CPublish *cut = new CPublish();
   class="type">uchar expected[] = {class="num">51, class="num">6, class="num">0, class="num">1, &class="macro">#x27;a&class="macro">#x27;, class="num">0, class="num">1, class="num">0}; class=class="str">"cmt">// QoS > class="num">0 require packet ID
   class="type">uchar result[];
   cut.SetTopicName("a");
   cut.SetRetain(true);
   cut.SetQoS_1(true);
   cut.Build(result);
   class="type">bool isTrue = AssertEqual(expected, result);
   class="kw">delete(cut);
   ZeroMemory(result);
   class="kw">return isTrue;
  }
class="type">void CPktPublish::SetRetain(const class="type">bool retain)
  {
   retain ? m_pubflags |= RETAIN_FLAG : m_pubflags &= ~RETAIN_FLAG;
  }
class="type">void CPktPublish::SetQoS_1(const class="type">bool QoS_1)
  {
   QoS_1 ? m_pubflags |= QoS_1_FLAG : m_pubflags &= ~QoS_1_FLAG;
  }
class="type">void CPktPublish::SetQoS_2(const class="type">bool QoS_2)
  {
   QoS_2 ? m_pubflags |= QoS_2_FLAG : m_pubflags &= ~QoS_2_FLAG;
  }
class="type">void CPktPublish::SetDup(const class="type">bool dup)
  {
   dup ? m_pubflags |= DUP_FLAG : m_pubflags &= ~DUP_FLAG;
  }
pkt[class="num">0] |= m_pubflags;
class="type">void CPktPublish::SetTopicName(const class="type">class="kw">string topic_name)
  {
   if(HasWildcardChar(topic_name) || StringLen(topic_name) == class="num">0)
     {
      ArrayFree(m_topname);
      class="kw">return;
     }
   EncodeUTF8String(topic_name, m_topname);
  }
class=class="str">"cmt">// QoS > class="num">0 requires packet ID
   if((m_pubflags & 0x06) != class="num">0)
     {
      SetPacketID(pkt, pkt.Size());
     }
class="macro">#define TEST true

报文ID与MQTT属性字段的落地写法

在 MT5 里手搓 MQTT 报文,先得解决可变头里的 Packet ID。上面那段 SetPacketID 用 MathSrand((int)TimeLocal()) 做随机种子,再 MathRand() 取一个 int,拆成高 8 位和低 8 位塞进 uchar 数组——MSB 右移 8 位,LSB 取模 256 并与 0xff。若数组扩容失败(ArrayResize 返回负值),直接 printf 报错并 return,不破坏原 buf。 回测或模拟环境里建议把 TEST 宏打开:代码会强制把 ID 写成 0x0001,并打印 WARN 提示。这样你在策略测试器里跑断点,能一眼区分真实随机包和固定测试包,避免把随机数带进历史回测造成不可复现。 MQTT 的 Properties 字段几乎挂在所有控制报文末尾(CONNECT / PUBLISH / SUBSCRIBE 等),Will Properties 里还能再嵌一组。协议用单字节 Identifier 区分类型:0x01 是载荷格式指示符(1 字节),0x02 是消息过期时间(4 字节整数),0x0B 是订阅标识符(变长字节整数),0x11 是会话过期间隔(4 字节整数)。 CPktPublish::SetPayloadFormatIndicator 展示了怎么拼一个属性:先建 2 字节 aux,aux[0] 放标识符 0x01,aux[1] 放枚举值,再用 ArrayCopy 接到 m_props 尾部。PAYLOAD_FORMAT_INDICATOR 里 RAW_BYTES=0x00 表示裸字节,UTF8=0x01 表示 UTF-8 字符串——接小布这类 AIGC 文本流时通常选后者。

MQL5 / C++
class="type">void SetPacketID(class="type">uchar& buf[], class="type">int start_idx)
  {
class=class="str">"cmt">// MathRand - Before the first call of the function, it&class="macro">#x27;s necessary to call
class=class="str">"cmt">// MathSrand to set the generator of pseudorandom numbers to the initial state.
   MathSrand((class="type">int)TimeLocal());
   class="type">int packet_id = MathRand();
   if(ArrayResize(buf, buf.Size() + class="num">2) < class="num">0)
     {
       printf("ERROR: failed to resize array at %s", __FUNCTION__);
       class="kw">return;
     }
   buf[start_idx] = (class="type">uchar)packet_id >> class="num">8; class=class="str">"cmt">// MSB
   buf[start_idx + class="num">1] = (class="type">uchar)(packet_id % class="num">256) & 0xff; class=class="str">"cmt">//LSB
class=class="str">"cmt">//--- if testing, set packet ID to class="num">1
   if(TEST)
     {
       Print("WARN: SetPacketID TEST true fixed ID = class="num">1");
       buf[start_idx] = class="num">0; class=class="str">"cmt">// MSB
       buf[start_idx + class="num">1] = class="num">1; class=class="str">"cmt">//LSB
     }
  }
class=class="str">"cmt">//+------------------------------------------------------------------+
class=class="str">"cmt">//|                    PROPERTIES                                    |
class=class="str">"cmt">//+------------------------------------------------------------------+
class=class="str">"cmt">/*
The last field in the Variable Header of the CONNECT, CONNACK, PUBLISH, PUBACK, PUBREC,
PUBREL, PUBCOMP, SUBSCRIBE, SUBACK, UNSUBSCRIBE, UNSUBACK, DISCONNECT, and
AUTH packet is a set of Properties. In the CONNECT packet there is also an optional set of Properties in
the Will Properties field with the Payload
*/
class="macro">#define MQTT_PROP_IDENTIFIER_PAYLOAD_FORMAT_INDICATOR          0x01 class=class="str">"cmt">// (class="num">1) Byte
class="macro">#define MQTT_PROP_IDENTIFIER_MESSAGE_EXPIRY_INTERVAL           0x02 class=class="str">"cmt">// (class="num">2) Four Byte Integer
class="macro">#define MQTT_PROP_IDENTIFIER_CONTENT_TYPE                      0x03 class=class="str">"cmt">// (class="num">3) UTF-class="num">8 Encoded String
class="macro">#define MQTT_PROP_IDENTIFIER_RESPONSE_TOPIC                    0x08 class=class="str">"cmt">// (class="num">8) UTF-class="num">8 Encoded String
class="macro">#define MQTT_PROP_IDENTIFIER_CORRELATION_DATA                  0x09 class=class="str">"cmt">// (class="num">9) Binary Data
class="macro">#define MQTT_PROP_IDENTIFIER_SUBSCRIPTION_IDENTIFIER           0x0B class=class="str">"cmt">// (class="num">11) Variable Byte Integer
class="macro">#define MQTT_PROP_IDENTIFIER_SESSION_EXPIRY_INTERVAL           0x11 class=class="str">"cmt">// (class="num">17) Four Byte Integer
class="type">void CPktPublish::SetPayloadFormatIndicator(PAYLOAD_FORMAT_INDICATOR format)
  {
   class="type">uchar aux[class="num">2];
   aux[class="num">0] = MQTT_PROP_IDENTIFIER_PAYLOAD_FORMAT_INDICATOR;
   aux[class="num">1] = (class="type">uchar)format;
   ArrayCopy(m_props, aux, m_props.Size());
  }
enum PAYLOAD_FORMAT_INDICATOR
  {
   RAW_BYTES   = 0x00,
   UTF8        = 0x01
  };

◍ Publish 包的属性字段怎么塞进字节流

MQTT 的 PUBLISH 报文靠属性标识符 + 编码值拼出 m_props 字节数组,上面这组方法把常用属性逐个追加进去。 消息过期间隔用 4 字节大端写入:先放标识符 MQTT_PROP_IDENTIFIER_MESSAGE_EXPIRY_INTERVAL(1 字节),再调 EncodeFourByteInteger 把 uint 拆成 4 个 uchar,高位在前。 主题别名只需 2 字节,EncodeTwoByteInteger 里 dest_buf[0] 存高 8 位、dest_buf[1] 存低 8 位,范围受 ushort 限制。 响应主题和相关数据走二进制拼接:响应主题先写标识符再 EncodeUTF8String;相关数据 binary_data 直接 ArrayCopy 进 m_props,不带长度头。 用户属性最啰嗦——标识符 + key 的 UTF8 串 + val 的 UTF8 串,连续三次 ArrayCopy,顺序不能反。 订阅标识符做了边界判断,subscript_id 不在 1 到 0xfffffff(268435455)之间就直接拦掉,避免变长编码越界。 在 MT5 里写 EA 做 MQTT 推送时,照这套顺序拼属性, broker 才可能正确解析;乱序或漏标识符会静默丢包。

MQL5 / C++
class="type">void CPktPublish::SetMessageExpiryInterval(class="type">uint msg_expiry_interval)
  {
   class="type">uchar aux[class="num">4];
   aux[class="num">0] = MQTT_PROP_IDENTIFIER_MESSAGE_EXPIRY_INTERVAL;
   ArrayCopy(m_props, aux, m_props.Size(), class="num">0, class="num">1);
   EncodeFourByteInteger(msg_expiry_interval, aux);
   ArrayCopy(m_props, aux, m_props.Size());
  }
class="type">void EncodeFourByteInteger(class="type">uint val, class="type">uchar &dest_buf[])
  {
   ArrayResize(dest_buf, class="num">4);
   dest_buf[class="num">0] = (class="type">uchar)(val >> class="num">24) & 0xff;
   dest_buf[class="num">1] = (class="type">uchar)(val >> class="num">16) & 0xff;
   dest_buf[class="num">2] = (class="type">uchar)(val >> class="num">8) & 0xff;
   dest_buf[class="num">3] = (class="type">uchar)val & 0xff;
  }
class="type">void CPktPublish::SetTopicAlias(class="type">class="kw">ushort topic_alias)
  {
   class="type">uchar aux[class="num">2];
   aux[class="num">0] = MQTT_PROP_IDENTIFIER_TOPIC_ALIAS;
   ArrayCopy(m_props, aux, m_props.Size(), class="num">0, class="num">1);
   EncodeTwoByteInteger(topic_alias, aux);
   ArrayCopy(m_props, aux, m_props.Size());
  }
class="type">void EncodeTwoByteInteger(class="type">uint val, class="type">uchar &dest_buf[])
  {
   ArrayResize(dest_buf, class="num">2);
   dest_buf[class="num">0] = (class="type">uchar)(val >> class="num">8) & 0xff;
   dest_buf[class="num">1] = (class="type">uchar)val & 0xff;
  }
class="type">void CPktPublish::SetResponseTopic(const class="type">class="kw">string response_topic)
  {
   class="type">uchar aux[class="num">1];
   aux[class="num">0] = MQTT_PROP_IDENTIFIER_RESPONSE_TOPIC;
   ArrayCopy(m_props, aux, m_props.Size());
   class="type">uchar buf[];
   EncodeUTF8String(response_topic, buf);
   ArrayCopy(m_props, buf, m_props.Size());
  }
class="type">void CPktPublish::SetCorrelationData(class="type">uchar &binary_data[])
  {
   class="type">uchar aux[class="num">1];
   aux[class="num">0] = MQTT_PROP_IDENTIFIER_CORRELATION_DATA;
   ArrayCopy(m_props, aux, m_props.Size());
   ArrayCopy(m_props, binary_data, m_props.Size());
  }
class="type">void CPktPublish::SetUserProperty(const class="type">class="kw">string key, const class="type">class="kw">string val)
  {
   class="type">uchar aux[class="num">1];
   aux[class="num">0] = MQTT_PROP_IDENTIFIER_USER_PROPERTY;
   ArrayCopy(m_props, aux, m_props.Size());
   class="type">uchar key_buf[];
   EncodeUTF8String(key, key_buf);
   ArrayCopy(m_props, key_buf, m_props.Size());
   class="type">uchar val_buf[];
   EncodeUTF8String(val, val_buf);
   ArrayCopy(m_props, val_buf, m_props.Size());
  }
class="type">void CPktPublish::SetSubscriptionIdentifier(class="type">uint subscript_id)
  {
   if(subscript_id < class="num">1 || subscript_id > 0xfffffff)
     {

「PUBLISH 报文拼装时的几个坑」

在 MT5 里手搓 MQTT 的 PUBLISH 包,最容易翻车的是订阅标识符范围。代码里硬性校验了 1 到 268,435,455,超出就 printf 报错直接 return,Broker 端大概率也会断连,外汇信号转发场景里别乱填 0。 SetContentType 和 SetPayload 都走 UTF-8 编码再 ArrayCopy 进属性或负载区,但注意 SetPayload 里用的是 m_props.Size() 作偏移而不是 m_payload 自身长度,初写容易把负载错位叠到属性尾上。 Build 函数才是真正封包的地方:pkt[0] 高 4 位是 PUBLISH 类型左移 4,低 4 位是发布标志;topic name 为空时直接报错返回,这是必检项。QoS>0 时靠 m_pubflags & 0x06 判断是否补包 ID,少了这步 Broker 会回 DISCONNECT。 剩余长度 m_remlen 在 Build 末尾才累加 pkt.Size()-2 并重新编码插到索引 1,意味着前面拼 topic、props、payload 时都还没算总长。开 MT5 把这段直接挂到 EA 的调试脚本里,用 PacketLogger 打印 pkt 数组,能立刻看到字节序对不对。

MQL5 / C++
   printf("Error: " + __FUNCTION__ +  "Subscription Identifier must be between class="num">1 and class="num">268,class="num">435,class="num">455");
    class="kw">return;
   }
  class="type">uchar aux[class="num">1];
  aux[class="num">0] = MQTT_PROP_IDENTIFIER_SUBSCRIPTION_IDENTIFIER;
  ArrayCopy(m_props, aux, m_props.Size());
  class="type">uchar buf[];
  EncodeVariableByteInteger(subscript_id, buf);
  ArrayCopy(m_props, buf, m_props.Size());
  }
class="type">void CPktPublish::SetContentType(const class="type">class="kw">string content_type)
  {
  class="type">uchar aux[class="num">1];
  aux[class="num">0] = MQTT_PROP_IDENTIFIER_CONTENT_TYPE;
  ArrayCopy(m_props, aux, m_props.Size());
  class="type">uchar buf[];
  EncodeUTF8String(content_type, buf);
  ArrayCopy(m_props, buf, m_props.Size());
  };
class="type">void CPktPublish::SetPayload(const class="type">class="kw">string payload)
  {
  class="type">uchar aux[];
  EncodeUTF8String(payload, aux);
  ArrayCopy(m_payload, aux, m_props.Size());
  }
class="type">void CPktPublish::Build(class="type">uchar &pkt[])
  {
  if(m_topname.Size() == class="num">0)
    {
      printf("Error: " + __FUNCTION__ + " topic name is mandatory");
      class="kw">return;
    }
  ArrayResize(pkt, class="num">2);
class=class="str">"cmt">// pkt type with publish flags
  pkt[class="num">0] = (class="type">uchar)PUBLISH << class="num">4;
  pkt[class="num">0] |= m_pubflags;
class=class="str">"cmt">// topic name
  ArrayCopy(pkt, m_topname, pkt.Size());
class=class="str">"cmt">// QoS > class="num">0 require packet ID
  if((m_pubflags & 0x06) != class="num">0)
    {
      SetPacketID(pkt, pkt.Size());
    }
class=class="str">"cmt">// properties length
  class="type">uchar buf[];
  EncodeVariableByteInteger(m_props.Size(), buf);
  ArrayCopy(pkt, buf, pkt.Size());
class=class="str">"cmt">// properties
  ArrayCopy(pkt, m_props, pkt.Size());
class=class="str">"cmt">// payload
  ArrayCopy(pkt, m_payload, pkt.Size());
class=class="str">"cmt">// remaining length
  m_remlen += pkt.Size() - class="num">2;
  class="type">uchar aux[];
  EncodeVariableByteInteger(m_remlen, aux);
  ArrayCopy(pkt, aux, class="num">1);
  }

PUBACK 确认逻辑与虚拟持久层

CPuback 类沿用了 CPublish 的骨架,同样挂在 IControlPacket 接口下。PUBACK 是 QoS 1 下对 PUBLISH 的应答包,固定报头两字节:首字节放控制包类型标识,次字节是剩余长度,位标志全保留。可变报头依次排着被确认包 ID、原因码、属性长度和属性。 协议是对称的,客户端和服务器都能当发送方或接收方,所以得补上接收者视角。没有 PUBACK,QoS 1 的 PUBLISH 就失去意义;而 PUBACK 自身又依赖持久层记住那些待确认的数据包 ID。眼下不必接真数据库,先弄个“像库一样”的桩函数 GetPendingPublishIDs 放在 DB.mqh 里,测试时返回几个写死的 ID 即可。 解码包 ID 要注意大端序:最高有效字节在前,用“高字节乘 256 加低字节”还原成 ushort。IsPendingPkt 先拉取待确认 ID 数组,再解出代理发来的 PUBACK 里的 ID 做比对,命中就返回 true,后续真实存储便可据此释放 ID。同一套思路稍后也能套到 QoS 2 的 PUBREC、PUBREL、PUBCOMP 上。 原因码位置固定,直接读可变报头第 3 字节(索引 2)就行。当前只在客户端接收端验证,未涉及发 PUBACK,上面几个函数足够跑通下一轮功能测试。待对接真代理时,字节序边界情况可能要补处理,但主逻辑不动。

MQL5 / C++
class="type">void GetPendingPublishIDs(class="type">class="kw">ushort &result[])
  {
   ArrayResize(result, class="num">3);
   result[class="num">0] = class="num">1;
   result[class="num">1] = class="num">255; class=class="str">"cmt">// one byte
   result[class="num">2] = class="num">65535; class=class="str">"cmt">// two bytes
  }
class="type">class="kw">ushort CPuback::GetPacketID(class="type">uchar &pkt[])
  {
   class="kw">return (pkt[class="num">0] * class="num">256) + pkt[class="num">1];
  }
class="type">bool CPuback::IsPendingPkt(class="type">uchar &pkt[])
{
   class="type">class="kw">ushort pending_ids[];
   GetPendingPublishIDs(pending_ids);
   class="type">class="kw">ushort packet_id = GetPacketID(pkt);
   for(class="type">uint i = class="num">0; i < pending_ids.Size(); i++)
     {
      if(pending_ids[i] == packet_id)
        {
         class="kw">return true;
        }
     }
   class="kw">return false;
}
class="type">uchar CPuback::GetReasonCode(class="type">uchar &pkt[])
  {
   class="kw">return pkt[class="num">2];
  }

◍ 把工具请下神坛

重构不是一次性任务,而是跟着功能长出来的习惯。MQTT 5.0 客户端的 CONNECT、CONNACK、PUBLISH、PUBACK 包跑通后,下一步才是 QoS 2 那套 PUBREC、PUBCOMP、PUBREL 对应包。 开源仓库里 MQTT.zip 是 20.83 KB、Tests.zip 是 16.88 KB,直接拖进 MT5 的 MQL5 目录就能编译验证。别把这套客户端当圣物,它只是你盯盘流程里一个可替换的管道。 真要落地,就开 MT5 把 Tests.zip 里的用例跑一遍,看 PUBACK 回包是否符合预期;不合就改类方法,改完再跑。工具归工具,盘面上的外汇和贵金属高风险不会因它消失,信号错了照样亏。

常见问题

按协议把 dup、qos、retain 三个位拼到首字节:qos=0 时 dup 必须为 0;retain 按业务需要置 1 或 0,别默认全开。
主题名紧跟固定头之后,用两字节长度前缀;报文ID 仅在 qos>0 时出现在主题名后面,两字节大端序。
可以,把拼装代码或字节流贴给小布,它会按协议核对标志位、主题长度和属性顺序,并指出常见越界点。
大概率是的。先用内存映射模拟持久层存报文ID,收到 PUBACK 再删;若仍丢,查发送端是否没等确认就清缓存。
最容易漏写属性长度前缀,或把字符串属性当成裸字节写;务必先写变长编码的长度,再写内容。