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;
|
|
}
|
|
};
|
|
//+------------------------------------------------------------------+
|