Article-24929-ZMQ-SUB/Zmtp.mqh

809 lines
27 KiB
MQL5

2026-09-23 12:31:36 +00:00
//+------------------------------------------------------------------+
//| Zmtp.mqh |
//| Copyright 2025, MetaQuotes Ltd. |
//| https://www.mql5.com |
//+------------------------------------------------------------------+
#property copyright "Copyright 2025, MetaQuotes Ltd."
#property link "https://www.mql5.com"
//+------------------------------------------------------------------+
//| Zmtp.mqh |
//| Minimal ZMTP 3.0 client transport for MetaTrader 5. |
//| |
//| Implements: |
//| - ZMTP 3.0 greeting + framing (RFC 23) |
//| - NULL and PLAIN security mechanisms (RFC 24) |
//| - CZmqReqSocket : REQ/REP pattern |
//| |
//| CURVE security and the other socket-pattern classes are NOT |
//| implemented yet |
//+------------------------------------------------------------------+
//+------------------------------------------------------------------+
//| Socket-Type property this library presents to the ZMTP peer |
//+------------------------------------------------------------------+
enum ENUM_ZMQ_SOCKET_TYPE
{
ZMQ_TYPE_REQ,
ZMQ_TYPE_REP,
ZMQ_TYPE_DEALER,
ZMQ_TYPE_ROUTER,
ZMQ_TYPE_PUB,
ZMQ_TYPE_SUB,
ZMQ_TYPE_PUSH,
ZMQ_TYPE_PULL,
ZMQ_TYPE_PAIR
};
enum ENUM_ZMTP_SECURITY
{
ZMTP_SEC_NULL, // no authentication, no encryption
ZMTP_SEC_PLAIN, // username/password, sent in clear text
ZMTP_SEC_CURVE // NOT IMPLEMENTED - see note at bottom of file
};
//+------------------------------------------------------------------+
//| |
//+------------------------------------------------------------------+
string ZmqSocketTypeName(ENUM_ZMQ_SOCKET_TYPE t)
{
switch(t)
{
case ZMQ_TYPE_REQ:
return "REQ";
case ZMQ_TYPE_REP:
return "REP";
case ZMQ_TYPE_DEALER:
return "DEALER";
case ZMQ_TYPE_ROUTER:
return "ROUTER";
case ZMQ_TYPE_PUB:
return "PUB";
case ZMQ_TYPE_SUB:
return "SUB";
case ZMQ_TYPE_PUSH:
return "PUSH";
case ZMQ_TYPE_PULL:
return "PULL";
case ZMQ_TYPE_PAIR:
return "PAIR";
}
return "";
}
//+------------------------------------------------------------------+
//| One ZMTP frame as read off / written to a node |
//+------------------------------------------------------------------+
struct ZmtpFrame
{
bool more; // MORE flag - another frame follows in this message
bool command; // COMMAND flag - control frame (handshake), not user data
uchar data[];
};
//+------------------------------------------------------------------+
//| Appends a ZMTP command-property (name/value pair) to a buffer. |
//| Wire format: name-len(1) name value-len(4,BE) value |
//+------------------------------------------------------------------+
void ZmtpAppendProperty(uchar &buf[], string name, const uchar &value[])
{
uchar name_bytes[];
StringToCharArray(name, name_bytes, 0, StringLen(name));
int nlen = ArraySize(name_bytes);
int vlen = ArraySize(value);
int old = ArraySize(buf);
ArrayResize(buf, old + 1 + nlen + 4 + vlen);
buf[old] = (uchar)nlen;
for(int i=0; i<nlen; i++)
buf[old+1+i] = name_bytes[i];
int off = old + 1 + nlen;
buf[off+0] = (uchar)((vlen>>24)&0xFF);
buf[off+1] = (uchar)((vlen>>16)&0xFF);
buf[off+2] = (uchar)((vlen>>8)&0xFF);
buf[off+3] = (uchar)(vlen&0xFF);
for(int i=0; i<vlen; i++)
buf[off+4+i] = value[i];
}
//+------------------------------------------------------------------+
//| Appends src onto the end of dst - used to reassemble a multipart |
//| ZMTP message into one flat byte buffer. |
//+------------------------------------------------------------------+
void ZmtpAppendBytes(uchar &dst[], const uchar &src[])
{
int old = ArraySize(dst);
int add = ArraySize(src);
ArrayResize(dst, old+add);
ArrayCopy(dst, src, old, 0, add);
}
//+------------------------------------------------------------------+
//| CZmtpTransport |
//| Low-level ZMTP 3.0 connection: raw socket, receive buffering, |
//| greeting exchange, security handshake, frame read/write. |
//| Every socket-pattern class (REQ, and later DEALER/PUB/SUB/etc.) |
//| is built on top of one of these. |
//+------------------------------------------------------------------+
class CZmtpTransport
{
private:
int m_socket;
uchar m_rxbuf[];
int m_rxlen; // valid bytes currently in m_rxbuf
int m_rxpos; // read cursor into m_rxbuf
int m_timeout_ms;
ENUM_ZMTP_SECURITY m_security;
string m_username;
string m_password;
bool EnsureBytes(int n);
void PopBytes(int n, uchar &out[]);
bool RawSend(const uchar &buf[], int len);
bool SendGreeting();
bool RecvGreeting(string &peer_mechanism);
bool DoPlainHandshake();
bool SendReady(ENUM_ZMQ_SOCKET_TYPE stype, string identity);
bool RecvReady();
public:
CZmtpTransport();
~CZmtpTransport();
bool Connect(string host, int port, ENUM_ZMTP_SECURITY sec, int timeout_ms=5000);
void Disconnect();
bool IsConnected();
void SetPlainCredentials(string user, string pass) { m_username=user; m_password=pass; }
// Full ZMTP handshake: greeting -> (security mechanism) -> READY/READY.
// Call once, immediately after Connect().
bool Handshake(ENUM_ZMQ_SOCKET_TYPE stype, string identity="");
bool SendFrame(const uchar &data[], bool more, bool command=false);
bool SendCommand(string name, const uchar &body[]);
bool RecvFrame(ZmtpFrame &frame);
};
//+------------------------------------------------------------------+
//| |
//+------------------------------------------------------------------+
CZmtpTransport::CZmtpTransport()
{
m_socket = INVALID_HANDLE;
m_rxlen = 0;
m_rxpos = 0;
m_timeout_ms = 5000;
m_security = ZMTP_SEC_NULL;
m_username = "";
m_password = "";
}
//+------------------------------------------------------------------+
//| |
//+------------------------------------------------------------------+
CZmtpTransport::~CZmtpTransport()
{
Disconnect();
}
//+------------------------------------------------------------------+
//| |
//+------------------------------------------------------------------+
bool CZmtpTransport::Connect(string host, int port, ENUM_ZMTP_SECURITY sec, int timeout_ms=5000)
{
m_security = sec;
m_timeout_ms = timeout_ms;
m_rxlen = 0;
m_rxpos = 0;
ArrayResize(m_rxbuf, 0);
m_socket = SocketCreate();
if(m_socket == INVALID_HANDLE)
{
Print(__FUNCTION__," : Zmtp: SocketCreate failed, error ", GetLastError());
return false;
}
if(!SocketConnect(m_socket, host, (uint)port, (uint)timeout_ms))
{
Print(__FUNCTION__," : Zmtp: SocketConnect to ", host, ":", port, " failed, error ", GetLastError());
SocketClose(m_socket);
m_socket = INVALID_HANDLE;
return false;
}
return true;
}
//+------------------------------------------------------------------+
//| |
//+------------------------------------------------------------------+
void CZmtpTransport::Disconnect()
{
if(m_socket != INVALID_HANDLE)
{
SocketClose(m_socket);
m_socket = INVALID_HANDLE;
}
}
//+------------------------------------------------------------------+
//| |
//+------------------------------------------------------------------+
bool CZmtpTransport::IsConnected()
{
return (m_socket != INVALID_HANDLE) && SocketIsConnected(m_socket);
}
//--- send every byte, looping past partial sends -------------------
bool CZmtpTransport::RawSend(const uchar &buf[], int len)
{
int sent_total = 0;
uchar chunk[];
uint start = GetTickCount();
while(sent_total < len)
{
int remaining = len - sent_total;
ArrayResize(chunk, remaining);
ArrayCopy(chunk, buf, 0, sent_total, remaining);
int sent = SocketSend(m_socket, chunk, remaining);
if(sent <= 0)
{
if(!SocketIsConnected(m_socket))
{
Print(__FUNCTION__," : Zmtp: connection dropped mid-send");
return false;
}
if((int)(GetTickCount()-start) > m_timeout_ms)
{
Print(__FUNCTION__," : Zmtp: send timed out");
return false;
}
Sleep(1);
continue;
}
sent_total += sent;
start = GetTickCount(); // reset timeout on forward progress
}
return true;
}
//--- block (up to m_timeout_ms) until n bytes are buffered ---------
bool CZmtpTransport::EnsureBytes(int n)
{
uint start = GetTickCount();
while((m_rxlen - m_rxpos) < n)
{
if(!SocketIsConnected(m_socket))
{
Print(__FUNCTION__," : Zmtp: connection dropped mid-receive");
return false;
}
uint avail = SocketIsReadable(m_socket);
if(avail > 0)
{
uchar tmp[];
ArrayResize(tmp, avail);
int got = SocketRead(m_socket, tmp, avail, 100);
if(got > 0)
{
int need = m_rxlen + got;
if(need > ArraySize(m_rxbuf))
ArrayResize(m_rxbuf, need);
for(int i=0; i<got; i++)
m_rxbuf[m_rxlen+i] = tmp[i];
m_rxlen += got;
start = GetTickCount(); // reset timeout on forward progress
}
}
else
Sleep(1);
/*if((int)(GetTickCount()-start) > m_timeout_ms)
{
Print(__FUNCTION__," : Zmtp: receive timed out waiting for ", n, " bytes");
return false;
}*/
}
return true;
}
//+------------------------------------------------------------------+
//| |
//+------------------------------------------------------------------+
void CZmtpTransport::PopBytes(int n, uchar &out[])
{
ArrayResize(out, n);
for(int i=0; i<n; i++)
out[i] = m_rxbuf[m_rxpos+i];
m_rxpos += n;
// compact the buffer once we've drained a meaningful chunk of it
if(m_rxpos > 4096)
{
int remain = m_rxlen - m_rxpos;
for(int i=0; i<remain; i++)
m_rxbuf[i] = m_rxbuf[m_rxpos+i];
m_rxlen = remain;
m_rxpos = 0;
}
}
//--- 64-byte ZMTP greeting -------------------------------------------
bool CZmtpTransport::SendGreeting()
{
uchar g[];
ArrayResize(g, 64);
ArrayFill(g, 0, 64, 0);
g[0] = 0xFF; // signature start
g[9] = 0x7F; // signature end
g[10] = 3; // version-major = 3 (ZMTP 3.x)
g[11] = 0; // version-minor
string mech = "NULL";
if(m_security == ZMTP_SEC_PLAIN)
mech = "PLAIN";
else
if(m_security == ZMTP_SEC_CURVE)
mech = "CURVE";
uchar mech_bytes[];
StringToCharArray(mech, mech_bytes, 0, StringLen(mech));
for(int i=0; i<ArraySize(mech_bytes) && i<20; i++)
g[12+i] = mech_bytes[i];
g[32] = 0; // as-server: this library only ever connects as a client
// g[33..63] filler, already zero
return RawSend(g, 64);
}
//+------------------------------------------------------------------+
//| |
//+------------------------------------------------------------------+
bool CZmtpTransport::RecvGreeting(string &peer_mechanism)
{
if(!EnsureBytes(64))
return false;
uchar g[];
PopBytes(64, g);
if(g[0] != 0xFF || g[9] != 0x7F)
{
Print(__FUNCTION__," : Zmtp: peer sent an invalid greeting signature - is this really a ZMTP 3.x endpoint?");
return false;
}
int major = g[10];
if(major < 3)
{
Print(__FUNCTION__," : Zmtp: peer speaks ZMTP ", major, ".x - this library only supports ZMTP 3.x");
return false;
}
uchar mech_bytes[20];
for(int i=0; i<20; i++)
mech_bytes[i] = g[12+i];
int mlen = 0;
while(mlen < 20 && mech_bytes[mlen] != 0)
mlen++;
peer_mechanism = CharArrayToString(mech_bytes, 0, mlen);
return true;
}
//--- frame I/O ---------------------------------------------------------
bool CZmtpTransport::SendFrame(const uchar &data[], bool more, bool command=false)
{
int len = ArraySize(data);
uchar flags = 0;
if(more)
flags |= 0x01;
bool long_form = (len > 255);
if(long_form)
flags |= 0x02;
if(command)
flags |= 0x04;
uchar header[];
if(!long_form)
{
ArrayResize(header, 2);
header[0] = flags;
header[1] = (uchar)len;
}
else
{
ArrayResize(header, 9);
header[0] = flags;
ulong ulen = (ulong)len;
for(int i=0; i<8; i++)
header[1+i] = (uchar)((ulen >> (8*(7-i))) & 0xFF);
}
int hlen = ArraySize(header);
uchar packet[];
ArrayResize(packet, hlen+len);
ArrayCopy(packet, header, 0, 0, hlen);
if(len > 0)
ArrayCopy(packet, data, hlen, 0, len);
return RawSend(packet, ArraySize(packet));
}
//+------------------------------------------------------------------+
//| |
//+------------------------------------------------------------------+
bool CZmtpTransport::RecvFrame(ZmtpFrame &frame)
{
if(!EnsureBytes(1))
return false;
uchar hdr1[];
PopBytes(1, hdr1);
uchar flags = hdr1[0];
frame.more = (flags & 0x01) != 0;
frame.command = (flags & 0x04) != 0;
bool long_form = (flags & 0x02) != 0;
ulong len = 0;
if(!long_form)
{
uchar l[];
if(!EnsureBytes(1))
return false;
PopBytes(1, l);
len = l[0];
}
else
{
uchar l[];
if(!EnsureBytes(8))
return false;
PopBytes(8, l);
for(int i=0; i<8; i++)
len = (len << 8) | l[i];
}
if(len > 0)
{
if(!EnsureBytes((int)len))
return false;
PopBytes((int)len, frame.data);
}
else
ArrayResize(frame.data, 0);
return true;
}
//+------------------------------------------------------------------+
//| |
//+------------------------------------------------------------------+
bool CZmtpTransport::SendCommand(string name, const uchar &body[])
{
uchar name_bytes[];
StringToCharArray(name, name_bytes, 0, StringLen(name));
int nlen = ArraySize(name_bytes);
uchar frame[];
ArrayResize(frame, 1+nlen+ArraySize(body));
frame[0] = (uchar)nlen;
for(int i=0; i<nlen; i++)
frame[1+i] = name_bytes[i];
for(int i=0; i<ArraySize(body); i++)
frame[1+nlen+i] = body[i];
return SendFrame(frame, false, true);
}
//--- PLAIN mechanism: HELLO(user,pass) -> WELCOME, then READY/READY --
bool CZmtpTransport::DoPlainHandshake()
{
uchar body[];
uchar uv[];
StringToCharArray(m_username, uv, 0, StringLen(m_username));
uchar pv[];
StringToCharArray(m_password, pv, 0, StringLen(m_password));
ZmtpAppendProperty(body, "Username", uv);
ZmtpAppendProperty(body, "Password", pv);
if(!SendCommand("HELLO", body))
return false;
ZmtpFrame f;
if(!RecvFrame(f))
return false;
if(!f.command)
{
Print(__FUNCTION__," : Zmtp/PLAIN: expected a WELCOME command, got a data frame");
return false;
}
int nlen = f.data[0];
string cname = CharArrayToString(f.data, 1, nlen);
if(cname == "ERROR")
{
Print(__FUNCTION__," : Zmtp/PLAIN: server rejected the supplied username/password");
return false;
}
if(cname != "WELCOME")
{
Print(__FUNCTION__," : Zmtp/PLAIN: expected WELCOME, got '", cname, "'");
return false;
}
return true;
}
//+------------------------------------------------------------------+
//| |
//+------------------------------------------------------------------+
bool CZmtpTransport::SendReady(ENUM_ZMQ_SOCKET_TYPE stype, string identity)
{
uchar body[];
ArrayResize(body, 0);
uchar stv[];
string sname = ZmqSocketTypeName(stype);
StringToCharArray(sname, stv, 0, StringLen(sname));
ZmtpAppendProperty(body, "Socket-Type", stv);
if(identity != "")
{
uchar idv[];
StringToCharArray(identity, idv, 0, StringLen(identity));
ZmtpAppendProperty(body, "Identity", idv);
}
return SendCommand("READY", body);
}
//+------------------------------------------------------------------+
//| |
//+------------------------------------------------------------------+
bool CZmtpTransport::RecvReady()
{
ZmtpFrame f;
if(!RecvFrame(f))
return false;
if(!f.command)
{
Print(__FUNCTION__," : Zmtp: expected a READY command, got a data frame");
return false;
}
int nlen = f.data[0];
string cname = CharArrayToString(f.data, 1, nlen);
if(cname == "ERROR")
{
Print(__FUNCTION__," : Zmtp: peer sent ERROR during handshake");
return false;
}
if(cname != "READY")
{
Print(__FUNCTION__," : Zmtp: expected READY, got '", cname, "'");
return false;
}
return true;
}
//+------------------------------------------------------------------+
//| |
//+------------------------------------------------------------------+
bool CZmtpTransport::Handshake(ENUM_ZMQ_SOCKET_TYPE stype, string identity="")
{
if(m_security == ZMTP_SEC_CURVE)
{
Print(__FUNCTION__," : Zmtp: ZMTP_SEC_CURVE is not implemented yet - see the note at the bottom of Zmtp.mqh");
return false;
}
if(!SendGreeting())
return false;
string peer_mech;
if(!RecvGreeting(peer_mech))
return false;
string want_mech = (m_security == ZMTP_SEC_PLAIN) ? "PLAIN" : "NULL";
if(peer_mech != want_mech)
{
Print(__FUNCTION__," : Zmtp: security mismatch - this socket requested ", want_mech,
" but the peer offered '", peer_mech, "'");
return false;
}
if(m_security == ZMTP_SEC_PLAIN)
if(!DoPlainHandshake())
return false;
if(!SendReady(stype, identity))
return false;
if(!RecvReady())
return false;
return true;
}
//+------------------------------------------------------------------+
//| CZmqReqSocket - ZMQ REQ pattern over ZMTP. |
//| |
//| Enforces the strict send/recv/send/recv alternation that REQ |
//| requires, and handles the empty delimiter-frame envelope that |
//| REQ/REP sockets use automatically in real ZeroMQ. |
//+------------------------------------------------------------------+
class CZmqReqSocket
{
private:
CZmtpTransport m_t;
bool m_awaiting_reply;
public:
CZmqReqSocket() { m_awaiting_reply = false; }
bool Connect(string host, int port, ENUM_ZMTP_SECURITY sec=ZMTP_SEC_NULL,
int timeout_ms=5000, string identity="", string plain_user="", string plain_pass="")
{
if(sec == ZMTP_SEC_PLAIN)
m_t.SetPlainCredentials(plain_user, plain_pass);
if(!m_t.Connect(host, port, sec, timeout_ms))
return false;
return m_t.Handshake(ZMQ_TYPE_REQ, identity);
}
void Disconnect() { m_t.Disconnect(); }
bool IsConnected() { return m_t.IsConnected(); }
// Single-part string request.
bool Send(const string &request)
{
uchar body[];
StringToCharArray(request, body, 0, StringLen(request));
return SendRaw(body);
}
// Single-part binary request (e.g. a packed struct of doubles).
bool SendRaw(const uchar &request[])
{
if(m_awaiting_reply)
{
Print(__FUNCTION__," : Zmq REQ: Send() called out of order - a reply is still pending");
return false;
}
uchar empty[];
ArrayResize(empty, 0);
// REQ envelope: empty delimiter frame, then the request body
if(!m_t.SendFrame(empty, true))
return false;
if(!m_t.SendFrame(request, false))
return false;
m_awaiting_reply = true;
return true;
}
// Blocks for the matching reply (up to the connect-time timeout).
bool Recv(string &reply)
{
uchar raw[];
if(!RecvRaw(raw))
return false;
reply = CharArrayToString(raw, 0, ArraySize(raw));
return true;
}
bool RecvRaw(uchar &reply[])
{
if(!m_awaiting_reply)
{
Print(__FUNCTION__," : Zmq REQ: Recv() called before Send()");
return false;
}
ZmtpFrame f;
if(!m_t.RecvFrame(f))
return false;
// If a ROUTER sits between us and the REP, extra identity/envelope
// frames may precede the empty delimiter - skip anything non-empty.
while(ArraySize(f.data) > 0 && f.more)
{
if(!m_t.RecvFrame(f))
return false;
}
ArrayResize(reply, 0);
while(f.more)
{
if(!m_t.RecvFrame(f))
return false;
int old = ArraySize(reply);
int add = ArraySize(f.data);
ArrayResize(reply, old+add);
ArrayCopy(reply, f.data, old, 0, add);
}
m_awaiting_reply = false;
return true;
}
};
//+------------------------------------------------------------------+
//| CZmqSubSocket - ZMQ SUB pattern over ZMTP. |
//| Receive-only, pairs with a PUB peer. Per ZMTP, subscriptions are |
//| sent to the peer as ordinary data frames (not COMMAND frames): |
//| a leading 0x01 (SUBSCRIBE) or 0x00 (UNSUBSCRIBE) byte followed by |
//| the topic-prefix bytes - filtering happens on the PUB side. |
//+------------------------------------------------------------------+
class CZmqSubSocket
{
private:
CZmtpTransport m_t;
bool SendSubscription(const string &prefix, uchar marker)
{
uchar p[];
StringToCharArray(prefix, p, 0, StringLen(prefix));
uchar frame[];
ArrayResize(frame, 1+ArraySize(p));
frame[0] = marker;
for(int i=0; i<ArraySize(p); i++)
frame[1+i] = p[i];
return m_t.SendFrame(frame, false);
}
public:
CZmqSubSocket(void) { }
bool Connect(string host, int port, ENUM_ZMTP_SECURITY sec=ZMTP_SEC_NULL,
int timeout_ms=5000, string identity="", string plain_user="", string plain_pass="")
{
if(sec == ZMTP_SEC_PLAIN)
m_t.SetPlainCredentials(plain_user, plain_pass);
if(!m_t.Connect(host, port, sec, timeout_ms))
return false;
return m_t.Handshake(ZMQ_TYPE_SUB, identity);
}
void Disconnect() { m_t.Disconnect(); }
bool IsConnected() { return m_t.IsConnected(); }
// Subscribe to any topic starting with `prefix` - pass "" to receive
// everything the publisher sends. Call before/after connecting; must
// be called at least once or a compliant PUB peer will send nothing.
bool Subscribe(const string &prefix) { return SendSubscription(prefix, 0x01); }
bool Unsubscribe(const string &prefix) { return SendSubscription(prefix, 0x00); }
// Blocks for the next published message, matching CZmqPubSocket::Publish's
// two-frame [topic][payload] convention.
bool Recv(string &topic, string &message)
{
ZmtpFrame f;
if(!m_t.RecvFrame(f))
return false;
topic = CharArrayToString(f.data, 0, ArraySize(f.data));
uchar raw[];
while(f.more)
{
if(!m_t.RecvFrame(f))
return false;
ZmtpAppendBytes(raw, f.data);
}
message = CharArrayToString(raw, 0, ArraySize(raw));
return true;
}
bool RecvRaw(uchar &topic[], uchar &message[])
{
ZmtpFrame f;
if(!m_t.RecvFrame(f))
return false;
ArrayResize(topic, ArraySize(f.data));
ArrayCopy(topic, f.data, 0, 0, ArraySize(f.data));
ArrayResize(message, 0);
while(f.more)
{
if(!m_t.RecvFrame(f))
return false;
ZmtpAppendBytes(message, f.data);
}
return true;
}
// For peers using the single-frame (topic-prefixed) publish convention
// instead of the two-frame one - returns the frame's raw content
// unsplit, so you can peel the prefix off yourself.
bool RecvSingleFrame(string &data)
{
uchar raw[];
ZmtpFrame f;
if(!m_t.RecvFrame(f))
return false;
ZmtpAppendBytes(raw, f.data);
while(f.more)
{
if(!m_t.RecvFrame(f))
return false;
ZmtpAppendBytes(raw, f.data);
}
data = CharArrayToString(raw, 0, ArraySize(raw));
return true;
}
};
//+------------------------------------------------------------------+