809 lines
27 KiB
MQL5
809 lines
27 KiB
MQL5
//+------------------------------------------------------------------+
| |||
//| 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;
| |||
}
| |||
};
| |||
//+------------------------------------------------------------------+
|