diff --git a/Include/Zmq/AtomicCounter.mqh b/Include/Zmq/AtomicCounter.mqh new file mode 100644 index 0000000..df57808 --- /dev/null +++ b/Include/Zmq/AtomicCounter.mqh @@ -0,0 +1,37 @@ +//+------------------------------------------------------------------+ +//| AtomicCounter.mqh | +//| Copyright 2016, Li Ding | +//| dingmaotu@hotmail.com | +//+------------------------------------------------------------------+ +#property copyright "Copyright 2016, Li Ding" +#property link "dingmaotu@hotmail.com" +#property strict + +#include "Common.mqh" + +#import "libzmq.dll" +intptr_t zmq_atomic_counter_new(void); +void zmq_atomic_counter_set(intptr_t counter,int value); +int zmq_atomic_counter_inc(intptr_t counter); +int zmq_atomic_counter_dec(intptr_t counter); +int zmq_atomic_counter_value(intptr_t counter); +void zmq_atomic_counter_destroy(intptr_t &counter_p); +#import +//+------------------------------------------------------------------+ +//| Atomic counter utility | +//+------------------------------------------------------------------+ +class AtomicCounter + { +private: + intptr_t m_ref; + +public: + AtomicCounter() {m_ref=zmq_atomic_counter_new();} + ~AtomicCounter() {zmq_atomic_counter_destroy(m_ref);} + + int increase() {return zmq_atomic_counter_inc(m_ref);} + int decrease() {return zmq_atomic_counter_dec(m_ref);} + int get() {return zmq_atomic_counter_value(m_ref);} + void set(int value) {zmq_atomic_counter_set(m_ref,value);} + }; +//+------------------------------------------------------------------+ diff --git a/Include/Zmq/Common.mqh b/Include/Zmq/Common.mqh new file mode 100644 index 0000000..650ad46 --- /dev/null +++ b/Include/Zmq/Common.mqh @@ -0,0 +1,116 @@ +//+------------------------------------------------------------------+ +//| Common.mqh | +//| Copyright 2016, Li Ding | +//| dingmaotu@hotmail.com | +//+------------------------------------------------------------------+ +#property copyright "Copyright 2016, Li Ding" +#property link "dingmaotu@hotmail.com" +#property strict + +#include "Errno.mqh" + +// Assume MT5 is 64bit, which is the default. +// Even though MT5 can be 32bit, there is no way to detect this +// by using preprocessor macros. Instead, MetaQuotes provides a +// function called IsX64 to detect this dynamically + +// This is just absurd. Why do you want to know the bitness of +// the runtime? To define pointer related entities at compile time! +// All integer types in MQL is uniform on both 32bit or 64bit +// architectures, so it is almost useless to have a runtime function IsX64. + +// Why not a __X64__? +#ifdef __MQL5__ +#define __X64__ +#endif + +#ifdef __X64__ +#define intptr_t long +#define uintptr_t ulong +#define size_t long +#else +#define intptr_t int +#define uintptr_t uint +#define size_t int +#endif + +#import "kernel32.dll" +void RtlMoveMemory(intptr_t dest,const uchar &array[],size_t length); +void RtlMoveMemory(uchar &array[],intptr_t src,size_t length); +int lstrlen(intptr_t psz); +int MultiByteToWideChar(uint codePage, + uint flags, + const intptr_t multiByteString, + int lengthMultiByte, + string &str, + int length + ); +#import +//+------------------------------------------------------------------+ +//| Copy the memory contents pointed by src to array | +//| array parameter should be initialized to the desired size | +//+------------------------------------------------------------------+ +void ArrayFromPointer(uchar &array[],intptr_t src,int count=WHOLE_ARRAY) + { + int size=(count==WHOLE_ARRAY)?ArraySize(array):count; + RtlMoveMemory(array,src,(size_t)size); + } +//+------------------------------------------------------------------+ +//| Copy array to the memory pointed by dest | +//+------------------------------------------------------------------+ +void ArrayToPointer(const uchar &array[],intptr_t dest,int count=WHOLE_ARRAY) + { + int size=(count==WHOLE_ARRAY)?ArraySize(array):count; + RtlMoveMemory(dest,array,(size_t)size); + } +//+------------------------------------------------------------------+ +//| Read a valid utf-8 string to the MQL environment | +//| With this function, there is no need to copy the string to char | +//| array, and convert with CharArrayToString | +//+------------------------------------------------------------------+ +string StringFromUtf8Pointer(intptr_t psz,int len) + { + if(len < 0) return NULL; + string res; + int required=MultiByteToWideChar(CP_UTF8,0,psz,len,res,0); + StringInit(res,required); + int resLength = MultiByteToWideChar(CP_UTF8,0,psz,len,res,required); + if(resLength != required) + { + return NULL; + } + else + { + return res; + } + } +//+------------------------------------------------------------------+ +//| for null-terminated string | +//+------------------------------------------------------------------+ +string StringFromUtf8Pointer(intptr_t psz) + { + int len=lstrlen(psz); + return StringFromUtf8Pointer(psz, len); + } +//+------------------------------------------------------------------+ +//| Convert a utf-8 byte array to a string | +//+------------------------------------------------------------------+ +string StringFromUtf8(const uchar &utf8[]) + { + return CharArrayToString(utf8, 0, -1, CP_UTF8); + } +//+------------------------------------------------------------------+ +//| Convert a string to a utf-8 byte array | +//+------------------------------------------------------------------+ +void StringToUtf8(const string str,uchar &utf8[],bool ending=true) + { + int count=ending ? -1 : StringLen(str); + StringToCharArray(str,utf8,0,count,CP_UTF8); + } + +#ifdef _DEBUG +#define Debug(msg) Print(">>> DEBUG: In ",__FUNCTION__,"(",__FILE__,":",__LINE__,") [", msg, "]") +#else +#define Debug(msg) +#endif +//+------------------------------------------------------------------+ diff --git a/Include/Zmq/Context.mqh b/Include/Zmq/Context.mqh new file mode 100644 index 0000000..587d903 --- /dev/null +++ b/Include/Zmq/Context.mqh @@ -0,0 +1,69 @@ +//+------------------------------------------------------------------+ +//| Context.mqh | +//| Copyright 2016, Li Ding | +//| dingmaotu@hotmail.com | +//+------------------------------------------------------------------+ + +#include "Common.mqh" +#include "SocketOptions.mqh" + +//--- Context options +#define ZMQ_IO_THREADS 1 +#define ZMQ_MAX_SOCKETS 2 +#define ZMQ_SOCKET_LIMIT 3 +#define ZMQ_THREAD_PRIORITY 3 +#define ZMQ_THREAD_SCHED_POLICY 4 +#define ZMQ_MAX_MSGSZ 5 + +//--- Default for new contexts +#define ZMQ_IO_THREADS_DFLT 1 +#define ZMQ_MAX_SOCKETS_DFLT 1023 +#define ZMQ_THREAD_PRIORITY_DFLT -1 +#define ZMQ_THREAD_SCHED_POLICY_DFLT -1 + +#import "libzmq.dll" +intptr_t zmq_ctx_new(void); +int zmq_ctx_term(intptr_t context); +int zmq_ctx_shutdown(intptr_t context); +int zmq_ctx_set(intptr_t context,int option,int optval); +int zmq_ctx_get(intptr_t context,int option); +#import +//+------------------------------------------------------------------+ +//| Wraps a 0MZ context | +//+------------------------------------------------------------------+ +class Context + { +private: + intptr_t m_ref; +protected: + int get(int option) {return zmq_ctx_get(m_ref,option);} + bool set(int option,int optval) {return 0==zmq_ctx_set(m_ref,option,optval);} +public: + Context() {m_ref=zmq_ctx_new();} + ~Context() {if(0!=zmq_ctx_term(m_ref)){Debug("failed to terminate context");}} + // for better cooperation between objects + intptr_t ref() const {return m_ref;} + bool shutdown() {return 0==zmq_ctx_shutdown(m_ref);} + + int getIoThreads() {return get(ZMQ_IO_THREADS);} + void setIoThreads(int value) {if(!set(ZMQ_IO_THREADS,value)){Debug("failed to set ZMQ_IO_THREADS");}} + + int getMaxSockets() {return get(ZMQ_MAX_SOCKETS);} + void setMaxSockets(int value) {if(!set(ZMQ_MAX_SOCKETS,value)){Debug("failed to set ZMQ_MAX_SOCKETS");}} + + int getMaxMessageSize() {return get(ZMQ_MAX_MSGSZ);} + void setMaxMessageSize(int value) {if(!set(ZMQ_MAX_MSGSZ,value)){Debug("failed to set ZMQ_MAX_MSGSZ");}} + + int getSocketLimit() {return get(ZMQ_SOCKET_LIMIT);} + + int getIpv6Options() {return get(ZMQ_IPV6);} + void setIpv6Options(int value) {if(!set(ZMQ_IPV6,value)){Debug("failed to set ZMQ_IPV6");}} + + bool isBlocky() {return 1==get(ZMQ_BLOCKY);} + void setBlocky(bool value) {if(!set(ZMQ_BLOCKY,value?1:0)){Debug("failed to set ZMQ_BLOCKY");}} + + //--- Following options is not supported on windows + void setSchedulingPolicy(int value) {/*ZMQ_THREAD_SCHED_POLICY*/} + void setThreadPriority(int value) {/*ZMQ_THREAD_PRIORITY*/} + }; +//+------------------------------------------------------------------+ diff --git a/Include/Zmq/Errno.mqh b/Include/Zmq/Errno.mqh new file mode 100644 index 0000000..000ece4 --- /dev/null +++ b/Include/Zmq/Errno.mqh @@ -0,0 +1,106 @@ +//+------------------------------------------------------------------+ +//| Errno.mqh | +//+------------------------------------------------------------------+ + +// Following error codes come from Microsoft CRT header errno.h + +// Error codes +#define EPERM 1 +#define ENOENT 2 +#define ESRCH 3 +#define EINTR 4 +#define EIO 5 +#define ENXIO 6 +#define E2BIG 7 +#define ENOEXEC 8 +#define EBADF 9 +#define ECHILD 10 +#define EAGAIN 11 +#define ENOMEM 12 +#define EACCES 13 +#define EFAULT 14 +#define EBUSY 16 +#define EEXIST 17 +#define EXDEV 18 +#define ENODEV 19 +#define ENOTDIR 20 +#define EISDIR 21 +#define ENFILE 23 +#define EMFILE 24 +#define ENOTTY 25 +#define EFBIG 27 +#define ENOSPC 28 +#define ESPIPE 29 +#define EROFS 30 +#define EMLINK 31 +#define EPIPE 32 +#define EDOM 33 +#define EDEADLK 36 +#define ENAMETOOLONG 38 +#define ENOLCK 39 +#define ENOSYS 40 +#define ENOTEMPTY 41 + +// Error codes used in the Secure CRT functions +#define EINVAL 22 +#define ERANGE 34 +#define EILSEQ 42 +#define STRUNCATE 80 + +// Support EDEADLOCK for compatibility with older Microsoft C versions +#define EDEADLOCK EDEADLK + +// POSIX Supplement +#define EALREADY 103 +#define EBADMSG 104 +#define ECANCELED 105 +#define EDESTADDRREQ 109 +#define EIDRM 111 +#define EISCONN 113 +#define ELOOP 114 +#define ENODATA 120 +#define ENOLINK 121 +#define ENOMSG 122 +#define ENOPROTOOPT 123 +#define ENOSR 124 +#define ENOSTR 125 +#define ENOTRECOVERABLE 127 +#define EOPNOTSUPP 130 +#define EOTHER 131 +#define EOVERFLOW 132 +#define EOWNERDEAD 133 +#define EPROTO 134 +#define EPROTOTYPE 136 +#define ETIME 137 +#define ETXTBSY 139 +#define EWOULDBLOCK 140 + +// Following error codes come from zmq.h +// 0MQ errors +#define ZMQ_HAUSNUMERO 156384712 + +#define ENOTSUP (ZMQ_HAUSNUMERO + 1) +#define EPROTONOSUPPORT (ZMQ_HAUSNUMERO + 2) +#define ENOBUFS (ZMQ_HAUSNUMERO + 3) +#define ENETDOWN (ZMQ_HAUSNUMERO + 4) +#define EADDRINUSE (ZMQ_HAUSNUMERO + 5) +#define EADDRNOTAVAIL (ZMQ_HAUSNUMERO + 6) +#define ECONNREFUSED (ZMQ_HAUSNUMERO + 7) +#define EINPROGRESS (ZMQ_HAUSNUMERO + 8) +#define ENOTSOCK (ZMQ_HAUSNUMERO + 9) +#define EMSGSIZE (ZMQ_HAUSNUMERO + 10) +#define EAFNOSUPPORT (ZMQ_HAUSNUMERO + 11) +#define ENETUNREACH (ZMQ_HAUSNUMERO + 12) +#define ECONNABORTED (ZMQ_HAUSNUMERO + 13) +#define ECONNRESET (ZMQ_HAUSNUMERO + 14) +#define ENOTCONN (ZMQ_HAUSNUMERO + 15) +#define ETIMEDOUT (ZMQ_HAUSNUMERO + 16) +#define EHOSTUNREACH (ZMQ_HAUSNUMERO + 17) +#define ENETRESET (ZMQ_HAUSNUMERO + 18) + +// Native 0MQ error codes +#define EFSM (ZMQ_HAUSNUMERO + 51) +#define ENOCOMPATPROTO (ZMQ_HAUSNUMERO + 52) +#define ETERM (ZMQ_HAUSNUMERO + 53) +#define EMTHREAD (ZMQ_HAUSNUMERO + 54) +//+------------------------------------------------------------------+ diff --git a/Include/Zmq/Socket.mqh b/Include/Zmq/Socket.mqh new file mode 100644 index 0000000..9ba51a7 --- /dev/null +++ b/Include/Zmq/Socket.mqh @@ -0,0 +1,293 @@ +//+------------------------------------------------------------------+ +//| Socket.mqh | +//| Copyright 2016, Li Ding | +//| dingmaotu@hotmail.com | +//+------------------------------------------------------------------+ +#property copyright "Copyright 2016, Li Ding" +#property link "dingmaotu@hotmail.com" +#property strict + +#include "Common.mqh" +#include "Context.mqh" +#include "SocketOptions.mqh" +#include "ZmqMsg.mqh" +//--- fd is SOCKET on Win32, which is defined as UINT_PTR +struct zmq_pollitem_t + { + intptr_t socket; + uintptr_t fd; + short events; + short revents; + }; + +//--- Socket types +#define ZMQ_PAIR 0 +#define ZMQ_PUB 1 +#define ZMQ_SUB 2 +#define ZMQ_REQ 3 +#define ZMQ_REP 4 +#define ZMQ_DEALER 5 +#define ZMQ_ROUTER 6 +#define ZMQ_PULL 7 +#define ZMQ_PUSH 8 +#define ZMQ_XPUB 9 +#define ZMQ_XSUB 10 +#define ZMQ_STREAM 11 + +//--- Message options +#define ZMQ_MORE 1 +#define ZMQ_SRCFD 2 // Deprecated +#define ZMQ_SHARED 3 + +//--- Send/recv options +#define ZMQ_DONTWAIT 1 +#define ZMQ_SNDMORE 2 + +//--- Socket transport events (TCP, IPC and TIPC only) +#define ZMQ_EVENT_CONNECTED 0x0001 +#define ZMQ_EVENT_CONNECT_DELAYED 0x0002 +#define ZMQ_EVENT_CONNECT_RETRIED 0x0004 +#define ZMQ_EVENT_LISTENING 0x0008 +#define ZMQ_EVENT_BIND_FAILED 0x0010 +#define ZMQ_EVENT_ACCEPTED 0x0020 +#define ZMQ_EVENT_ACCEPT_FAILED 0x0040 +#define ZMQ_EVENT_CLOSED 0x0080 +#define ZMQ_EVENT_CLOSE_FAILED 0x0100 +#define ZMQ_EVENT_DISCONNECTED 0x0200 +#define ZMQ_EVENT_MONITOR_STOPPED 0x0400 +#define ZMQ_EVENT_ALL 0xFFFF + +//--- I/O multiplexing +#define ZMQ_POLLIN 1 +#define ZMQ_POLLOUT 2 +// We only use 0MQ sockets on Windows +// So POLLERR and POLLPRI is of no use +#define ZMQ_POLLERR 4 +#define ZMQ_POLLPRI 8 + +#define ZMQ_POLLITEMS_DFLT 16 + +#import "libzmq.dll" +//+------------------------------------------------------------------+ +//| Sockets | +//+------------------------------------------------------------------+ +intptr_t zmq_socket(intptr_t context,int type); +int zmq_close(intptr_t s); +int zmq_bind(intptr_t s,const char &addr[]); +int zmq_connect(intptr_t s,const char &addr[]); +int zmq_unbind(intptr_t s,const char &addr[]); +int zmq_disconnect(intptr_t s,const char &addr[]); +int zmq_send(intptr_t s,const uchar &buf[],size_t len,int flags); +int zmq_send_const(intptr_t s,const uchar &buf[],size_t len,int flags); +int zmq_recv(intptr_t s,uchar &buf[],size_t len,int flags); +int zmq_socket_monitor(intptr_t s,const char &addr[],int events); +//+------------------------------------------------------------------+ +//| Message | +//+------------------------------------------------------------------+ +int zmq_msg_send(zmq_msg_t &msg,intptr_t s,int flags); +int zmq_msg_recv(zmq_msg_t &msg,intptr_t s,int flags); +//+------------------------------------------------------------------+ +//| I/O multiplexing | +//+------------------------------------------------------------------+ +int zmq_poll(zmq_pollitem_t &items[],int nitems,long timeout); +//+------------------------------------------------------------------+ +//| Message proxying | +//+------------------------------------------------------------------+ +int zmq_proxy(intptr_t frontend_ref,intptr_t backend_ref,intptr_t capture_ref); +int zmq_proxy_steerable(intptr_t frontend_ref,intptr_t backend_ref,intptr_t capture_ref,intptr_t control_ref); +#import +//+------------------------------------------------------------------+ +//| Wraps a 0MQ socket | +//+------------------------------------------------------------------+ +class Socket: public SocketOptions + { +public: + //--- it is not recommended to use this constructor directly: use Context factory methods instead + Socket(const Context &ctx,int type):SocketOptions(zmq_socket(ctx.ref(),type)){} + virtual ~Socket() {if(0!=zmq_close(m_ref)){Debug(StringFormat("Failed to close socket 0x%0X",m_ref));}} + + // for better cooperation between objects + intptr_t ref() const {return m_ref;} + + bool valid() const {return m_ref!=0;} + + //--- see Zmq::error() if any of the following command failed + //--- connection management + bool bind(string addr); + bool unbind(string addr); + bool connect(string addr); + bool disconnect(string addr); + + //--- send and receive packets + bool recv(uchar &buf[],bool nowait=false); + bool send(const uchar &buf[],bool nowait=false,bool more=false); + bool sendConst(const uchar &buf[],bool nowait=false,bool more=false); + + bool send(ZmqMsg &msg,bool nowait=false,bool more=false); + bool recv(ZmqMsg &msg,bool nowait=false); + + void register(zmq_pollitem_t &pollitem,bool read=false,bool write=false); + void register(zmq_pollitem_t &pollitems[],int index,bool read=false,bool write=false); + + //--- monitor socket events + bool monitor(string addr,int events); + + //--- proxy + static bool proxy(Socket *frontend,Socket *backend,Socket *capture); + static bool proxySteerable(Socket *frontend,Socket *backend,Socket *capture,Socket *control); + + //--- poll + static int poll(zmq_pollitem_t &arr[],long timeout); + }; +//+------------------------------------------------------------------+ +//| | +//+------------------------------------------------------------------+ +bool Socket::bind(string addr) + { + char arr[]; + StringToUtf8(addr,arr); + bool res=(0==zmq_bind(m_ref,arr)); + ArrayFree(arr); + return res; + } +//+------------------------------------------------------------------+ +//| | +//+------------------------------------------------------------------+ +bool Socket::unbind(string addr) + { + char arr[]; + StringToUtf8(addr,arr); + bool res=(0==zmq_unbind(m_ref,arr)); + ArrayFree(arr); + return res; + } +//+------------------------------------------------------------------+ +//| | +//+------------------------------------------------------------------+ +bool Socket::connect(string addr) + { + char arr[]; + StringToUtf8(addr,arr); + bool res=(0==zmq_connect(m_ref,arr)); + ArrayFree(arr); + return res; + } +//+------------------------------------------------------------------+ +//| | +//+------------------------------------------------------------------+ +bool Socket::disconnect(string addr) + { + char arr[]; + StringToUtf8(addr,arr); + bool res=(0==zmq_disconnect(m_ref,arr)); + ArrayFree(arr); + return res; + } +//+------------------------------------------------------------------+ +//| | +//+------------------------------------------------------------------+ +bool Socket::recv(uchar &buf[],bool nowait=false) + { + int options=0; + if(nowait) options|=ZMQ_DONTWAIT; + return 0==zmq_recv(m_ref,buf,ArraySize(buf),options); + } +//+------------------------------------------------------------------+ +//| | +//+------------------------------------------------------------------+ +bool Socket::send(const uchar &buf[],bool nowait=false,bool more=false) + { + int options=0; + if(nowait) options|=ZMQ_DONTWAIT; + if(more) options|=ZMQ_SNDMORE; + return 0==zmq_send(m_ref,buf,ArraySize(buf),options); + } +//+------------------------------------------------------------------+ +//| | +//+------------------------------------------------------------------+ +bool Socket::sendConst(const uchar &buf[],bool nowait=false,bool more=false) + { + int options=0; + if(nowait) options|=ZMQ_DONTWAIT; + if(more) options|=ZMQ_SNDMORE; + return 0==zmq_send_const(m_ref,buf,ArraySize(buf),options); + } +//+------------------------------------------------------------------+ +//| Send a zmq_msg_t through a socket | +//+------------------------------------------------------------------+ +bool Socket::send(ZmqMsg &msg,bool nowait=false,bool more=false) + { + int flags=0; + if(nowait) flags|=ZMQ_DONTWAIT; + if(more) flags|=ZMQ_SNDMORE; + return -1!=zmq_msg_send(msg,m_ref,flags); + } +//+------------------------------------------------------------------+ +//| Receive a zmq_msg_t from a socket | +//+------------------------------------------------------------------+ +bool Socket::recv(ZmqMsg &msg,bool nowait=false) + { + int flags=0; + if(nowait) flags|=ZMQ_DONTWAIT; + return -1!=zmq_msg_recv(msg,m_ref,flags); + } +//+------------------------------------------------------------------+ +//| | +//+------------------------------------------------------------------+ +bool Socket::monitor(string addr,int events) + { + uchar str[]; + StringToUtf8(addr,str); + bool res=(0==zmq_socket_monitor(m_ref,str,events)); + ArrayFree(str); + return res; + } +//+------------------------------------------------------------------+ +//| | +//+------------------------------------------------------------------+ +void Socket::register(zmq_pollitem_t &pollitem,bool read=false,bool write=false) + { + ZeroMemory(pollitem); + pollitem.socket=m_ref; + if(read) pollitem.events|=ZMQ_POLLIN; + if(write) pollitem.events|=ZMQ_POLLOUT; + } +//+------------------------------------------------------------------+ +//| | +//+------------------------------------------------------------------+ +void Socket::register(zmq_pollitem_t &pollitems[],int index,bool read=false,bool write=false) + { + ZeroMemory(pollitems[index]); + pollitems[index].socket=m_ref; + if(read) pollitems[index].events|=ZMQ_POLLIN; + if(write) pollitems[index].events|=ZMQ_POLLOUT; + } +//+------------------------------------------------------------------+ +//| | +//+------------------------------------------------------------------+ +bool Socket::proxy(Socket *frontend,Socket *backend,Socket *capture) + { + intptr_t frontend_ref= CheckPointer(frontend)==POINTER_DYNAMIC?frontend.ref():0; + intptr_t backend_ref = CheckPointer(backend)==POINTER_DYNAMIC?backend.ref():0; + intptr_t capture_ref=CheckPointer(capture)==POINTER_DYNAMIC?capture.ref():0; + return 0==zmq_proxy(frontend_ref, backend_ref, capture_ref); + } +//+------------------------------------------------------------------+ +//| | +//+------------------------------------------------------------------+ +bool Socket::proxySteerable(Socket *frontend,Socket *backend,Socket *capture,Socket *control) + { + intptr_t frontend_ref= CheckPointer(frontend)==POINTER_DYNAMIC?frontend.ref():0; + intptr_t backend_ref = CheckPointer(backend)==POINTER_DYNAMIC?backend.ref():0; + intptr_t capture_ref=CheckPointer(capture)==POINTER_DYNAMIC?capture.ref():0; + intptr_t control_ref=CheckPointer(control)==POINTER_DYNAMIC?control.ref():0; + return 0==zmq_proxy_steerable(frontend_ref, backend_ref, capture_ref, control_ref); + } +//+------------------------------------------------------------------+ +//| | +//+------------------------------------------------------------------+ +int Socket::poll(zmq_pollitem_t &arr[],long timeout) + { + return zmq_poll(arr,ArraySize(arr),timeout); + } +//+------------------------------------------------------------------+ diff --git a/Include/Zmq/SocketOptions.mqh b/Include/Zmq/SocketOptions.mqh new file mode 100644 index 0000000..bbf602f --- /dev/null +++ b/Include/Zmq/SocketOptions.mqh @@ -0,0 +1,342 @@ +//+------------------------------------------------------------------+ +//| Socket.mqh | +//| Copyright 2016, Li Ding | +//| dingmaotu@hotmail.com | +//+------------------------------------------------------------------+ +#property copyright "Copyright 2016, Li Ding" +#property link "dingmaotu@hotmail.com" +#property strict + +#include "Common.mqh" + +#import "libzmq.dll" +// We can overload the same function for different data types +// as in the C level the optval paramter is just a pointer +#define SOCKOPT_OVERLOAD_ARRAY(TYPE) \ +int zmq_setsockopt(intptr_t s,int option,const TYPE &optval[],\ + size_t optvallen);\ +int zmq_getsockopt(intptr_t s,int option,TYPE &optval[],\ + size_t &optvallen);\ + +#define SOCKOPT_OVERLOAD(TYPE) \ +int zmq_setsockopt(intptr_t s,int option,const TYPE &optval,\ + size_t optvallen);\ +int zmq_getsockopt(intptr_t s,int option,TYPE &optval,\ + size_t &optvallen);\ + +SOCKOPT_OVERLOAD_ARRAY(uchar) +SOCKOPT_OVERLOAD(long) +SOCKOPT_OVERLOAD(ulong) +SOCKOPT_OVERLOAD(int) +SOCKOPT_OVERLOAD(uint) +#import + +// Socket options +#define ZMQ_AFFINITY 4 +#define ZMQ_IDENTITY 5 +#define ZMQ_SUBSCRIBE 6 +#define ZMQ_UNSUBSCRIBE 7 +#define ZMQ_RATE 8 +#define ZMQ_RECOVERY_IVL 9 +#define ZMQ_SNDBUF 11 +#define ZMQ_RCVBUF 12 +#define ZMQ_RCVMORE 13 +#define ZMQ_FD 14 +#define ZMQ_EVENTS 15 +#define ZMQ_TYPE 16 +#define ZMQ_LINGER 17 +#define ZMQ_RECONNECT_IVL 18 +#define ZMQ_BACKLOG 19 +#define ZMQ_RECONNECT_IVL_MAX 21 +#define ZMQ_MAXMSGSIZE 22 +#define ZMQ_SNDHWM 23 +#define ZMQ_RCVHWM 24 +#define ZMQ_MULTICAST_HOPS 25 +#define ZMQ_RCVTIMEO 27 +#define ZMQ_SNDTIMEO 28 +#define ZMQ_LAST_ENDPOINT 32 +#define ZMQ_ROUTER_MANDATORY 33 +#define ZMQ_TCP_KEEPALIVE 34 +#define ZMQ_TCP_KEEPALIVE_CNT 35 +#define ZMQ_TCP_KEEPALIVE_IDLE 36 +#define ZMQ_TCP_KEEPALIVE_INTVL 37 +#define ZMQ_IMMEDIATE 39 +#define ZMQ_XPUB_VERBOSE 40 +#define ZMQ_ROUTER_RAW 41 +#define ZMQ_IPV6 42 +#define ZMQ_MECHANISM 43 +#define ZMQ_PLAIN_SERVER 44 +#define ZMQ_PLAIN_USERNAME 45 +#define ZMQ_PLAIN_PASSWORD 46 +#define ZMQ_CURVE_SERVER 47 +#define ZMQ_CURVE_PUBLICKEY 48 +#define ZMQ_CURVE_SECRETKEY 49 +#define ZMQ_CURVE_SERVERKEY 50 +#define ZMQ_PROBE_ROUTER 51 +#define ZMQ_REQ_CORRELATE 52 +#define ZMQ_REQ_RELAXED 53 +#define ZMQ_CONFLATE 54 +#define ZMQ_ZAP_DOMAIN 55 +#define ZMQ_ROUTER_HANDOVER 56 +#define ZMQ_TOS 57 +#define ZMQ_CONNECT_RID 61 +#define ZMQ_GSSAPI_SERVER 62 +#define ZMQ_GSSAPI_PRINCIPAL 63 +#define ZMQ_GSSAPI_SERVICE_PRINCIPAL 64 +#define ZMQ_GSSAPI_PLAINTEXT 65 +#define ZMQ_HANDSHAKE_IVL 66 +#define ZMQ_SOCKS_PROXY 68 +#define ZMQ_XPUB_NODROP 69 +#define ZMQ_BLOCKY 70 +#define ZMQ_XPUB_MANUAL 71 +#define ZMQ_XPUB_WELCOME_MSG 72 +#define ZMQ_STREAM_NOTIFY 73 +#define ZMQ_INVERT_MATCHING 74 +#define ZMQ_HEARTBEAT_IVL 75 +#define ZMQ_HEARTBEAT_TTL 76 +#define ZMQ_HEARTBEAT_TIMEOUT 77 +#define ZMQ_XPUB_VERBOSER 78 +#define ZMQ_CONNECT_TIMEOUT 79 +#define ZMQ_TCP_MAXRT 80 +#define ZMQ_THREAD_SAFE 81 +#define ZMQ_MULTICAST_MAXTPDU 84 +#define ZMQ_VMCI_BUFFER_SIZE 85 +#define ZMQ_VMCI_BUFFER_MIN_SIZE 86 +#define ZMQ_VMCI_BUFFER_MAX_SIZE 87 +#define ZMQ_VMCI_CONNECT_TIMEOUT 88 +#define ZMQ_USE_FD 89 +//+------------------------------------------------------------------+ +//| A dedicated class to get/set socket options | +//+------------------------------------------------------------------+ +class SocketOptions + { +protected: + intptr_t m_ref; + +#define SOCKOPT_WRAP_ARRAY(TYPE) \ + bool setOption(int option,const TYPE &value[],size_t len) {return 0==zmq_setsockopt(m_ref,option,value,len);}\ + bool getOption(int option,TYPE &value[],size_t &len) {return 0==zmq_getsockopt(m_ref,option,value,len);} + +#define SOCKOPT_WRAP(TYPE) \ + bool setOption(int option,TYPE value) {return 0==zmq_setsockopt(m_ref,option,value,sizeof(TYPE));}\ + bool getOption(int option,TYPE &value) {size_t s=sizeof(TYPE); return 0==zmq_getsockopt(m_ref,option,value,s);} + + SOCKOPT_WRAP_ARRAY(uchar) + SOCKOPT_WRAP(int) + SOCKOPT_WRAP(uint) + SOCKOPT_WRAP(long) + SOCKOPT_WRAP(ulong) + + bool setStringOption(int option,string value,bool ending=true); + bool getStringOption(int option,string &value,size_t length=1024); + SocketOptions(intptr_t ref):m_ref(ref){} +public: + //--- option templates + //--- various integer options +#define SOCKOPT_GET(TYPE, NAME, MACRO) \ + bool get##NAME(TYPE &value) {return getOption(MACRO,value);} +#define SOCKOPT_SET(TYPE, NAME, MACRO) \ + bool set##NAME(TYPE value) {return setOption(MACRO,value);} +#define SOCKOPT(TYPE,NAME,MACRO) \ + SOCKOPT_GET(TYPE,NAME,MACRO) \ + SOCKOPT_SET(TYPE,NAME,MACRO) + + //--- boolean options +#define SOCKOPT_GET_BOOL(NAME, MACRO) \ + bool is##NAME(bool &value) {int v; bool res=getOption(MACRO,v); value=(v==1);return res;} +#define SOCKOPT_SET_BOOL(NAME, MACRO) \ + bool set##NAME(bool value) {return setOption(MACRO,value?1:0);} +#define SOCKOPT_BOOL(NAME,MACRO) \ + SOCKOPT_GET_BOOL(NAME,MACRO) \ + SOCKOPT_SET_BOOL(NAME,MACRO) + + //--- null-terminated string options +#define SOCKOPT_GET_NTSTR(NAME,MACRO) \ + bool get##NAME(string &value) {return getStringOption(MACRO,value);} +#define SOCKOPT_SET_NTSTR(NAME,MACRO) \ + bool set##NAME(string value) {return setStringOption(MACRO,value);} +#define SOCKOPT_NTSTR(NAME,MACRO) \ + SOCKOPT_GET_NTSTR(NAME,MACRO) \ + SOCKOPT_SET_NTSTR(NAME,MACRO) + + //--- bytes array or string converted options +#define SOCKOPT_SET_BYTES(OptionName,Macro) \ + bool set##OptionName(const uchar &value[]) {return setOption(Macro,value,(size_t)ArraySize(value));} \ + bool set##OptionName(string value) {return setStringOption(Macro,value,false);} +#define SOCKOPT_GET_BYTES(OptionName,Macro,InitSize) \ + bool get##OptionName(uchar &value[]) {size_t len=(size_t)InitSize; ArrayResize(value,(int)len); bool res=getOption(Macro,value,len); if(res){ArrayResize(value,(int)len);}return res;} \ + bool get##OptionName(string &value) {return getStringOption(Macro,value,InitSize);} +#define SOCKOPT_BYTES(OptionName,Macro,InitSize) \ + SOCKOPT_SET_BYTES(OptionName,Macro) \ + SOCKOPT_GET_BYTES(OptionName,Macro,InitSize) + + //--- for curve key +#define SOCKOPT_CURVE_KEY(KeyType,Macro) \ + bool getCurve##KeyType##Key(uchar &key[32]) {size_t len=32; return getOption(Macro,key,len);} \ + bool getCurve##KeyType##Key(string &key) {return getStringOption(Macro,key,41);} \ + bool setCurve##KeyType##Key(const uchar &key[32]) {return setOption(Macro,key,32);} \ + bool setCurve##KeyType##Key(string key) {return setStringOption(Macro,key);} + + SOCKOPT_GET(int,Type,ZMQ_TYPE) + SOCKOPT(ulong,Affinity,ZMQ_AFFINITY) //64bit bitmask + SOCKOPT(int,BackLog,ZMQ_BACKLOG) // number of connections + SOCKOPT(int,Timeout,ZMQ_CONNECT_TIMEOUT) // milliseconds + SOCKOPT_GET_BOOL(ThreadSafe,ZMQ_THREAD_SAFE) + + SOCKOPT_SET_BOOL(Conflate,ZMQ_CONFLATE) // only for ZMQ_PULL, ZMQ_PUSH, ZMQ_SUB, ZMQ_PUB, ZMQ_DEALER types + + SOCKOPT_GET(int,Events,ZMQ_EVENTS); // bitmask of ZMQ_POLLIN, ZMQ_POLLOUT + + SOCKOPT_GET(uintptr_t,FileDescriptor,ZMQ_FD) + + SOCKOPT_GET(int,Mechanism,ZMQ_MECHANISM) // current security mechanism + + //--- plain + SOCKOPT_NTSTR(PlainUsername,ZMQ_PLAIN_USERNAME) + SOCKOPT_NTSTR(PlainPassword,ZMQ_PLAIN_PASSWORD) + SOCKOPT_BOOL(PlainServer,ZMQ_PLAIN_SERVER) + + //--- gssapi: this is not supported in this binding. Following methods will FAIL if you invoke them + SOCKOPT_BOOL(GssApiPlainText,ZMQ_GSSAPI_PLAINTEXT) + SOCKOPT_BOOL(GssApiServer,ZMQ_GSSAPI_SERVER) + SOCKOPT_NTSTR(GssApiPrincipal,ZMQ_GSSAPI_PRINCIPAL) + SOCKOPT_NTSTR(GssApiServicePrincipal,ZMQ_GSSAPI_SERVICE_PRINCIPAL) + + //--- curve + SOCKOPT_CURVE_KEY(Public,ZMQ_CURVE_PUBLICKEY) + SOCKOPT_CURVE_KEY(Secret,ZMQ_CURVE_SECRETKEY) + SOCKOPT_CURVE_KEY(Server,ZMQ_CURVE_SERVERKEY) + + SOCKOPT_SET_BOOL(CurveServer,ZMQ_CURVE_SERVER) + + SOCKOPT_GET_NTSTR(LastEndpoint,ZMQ_LAST_ENDPOINT) + + SOCKOPT(int,HandshakeInterval,ZMQ_HANDSHAKE_IVL) // milliseconds + SOCKOPT_SET(int,HeartbeatInterval,ZMQ_HEARTBEAT_IVL) // milliseconds + SOCKOPT_SET(int,HeartbeatTimeout,ZMQ_HEARTBEAT_TIMEOUT) // milliseconds + SOCKOPT_SET(int,HeartbeatTTL,ZMQ_HEARTBEAT_TTL) // milliseconds + + SOCKOPT_BOOL(Immediate,ZMQ_IMMEDIATE) + SOCKOPT_BOOL(Ipv6,ZMQ_IPV6) //--- ZMQ_IPV4ONLY is deprecated, use this instead + SOCKOPT(int,Linger,ZMQ_LINGER) // milliseconds + SOCKOPT(long,MaxMessageSize,ZMQ_MAXMSGSIZE) + + //--- multicast + SOCKOPT(int,MulticastHops,ZMQ_MULTICAST_HOPS) // hops + SOCKOPT(int,MulticastMaxTPDU,ZMQ_MULTICAST_MAXTPDU) // bytes + SOCKOPT(int,MulticastRate,ZMQ_RATE) // kilobits per second + SOCKOPT(int,RecoveryInterval,ZMQ_RECOVERY_IVL) // multicast recovery interval + + //--- there is a problem here: FileDescriptor should be SOCKET type on Windows, + //--- while in the zmq doc it is a int. Possible error in the doc + SOCKOPT(uintptr_t,UseFileDescriptor,ZMQ_USE_FD) + + SOCKOPT_SET_BOOL(ProbeRouter,ZMQ_PROBE_ROUTER) // only for ZMQ_ROUTER, ZMQ_DEALER, ZMQ_REQ + + SOCKOPT(int,ReceiveBuffer,ZMQ_RCVBUF) // bytes + SOCKOPT(int,ReceiveHighWaterMark,ZMQ_RCVHWM) // messages + SOCKOPT(int,ReceiveTimeout,ZMQ_RCVTIMEO) // milliseconds + SOCKOPT(int,SendBuffer,ZMQ_SNDBUF) // bytes + SOCKOPT(int,SendHighWaterMark,ZMQ_SNDHWM) // messages + SOCKOPT(int,SendTimout,ZMQ_SNDTIMEO) + + SOCKOPT_GET_BOOL(ReceiveMore,ZMQ_RCVMORE) + + SOCKOPT(int,ReconnectInterval,ZMQ_RECONNECT_IVL) // milliseconds + SOCKOPT(int,ReconnectIntervalMax,ZMQ_RECONNECT_IVL_MAX) // milliseconds + + //--- only for ZMQ_REQ + SOCKOPT_SET_BOOL(RequestCorrelated,ZMQ_REQ_CORRELATE) + SOCKOPT_SET_BOOL(RequestRelaxed,ZMQ_REQ_RELAXED) + + //--- only for ZMQ_SUB + SOCKOPT_SET_BYTES(Subscribe,ZMQ_SUBSCRIBE) + SOCKOPT_SET_BYTES(Unsubscribe,ZMQ_UNSUBSCRIBE) + + //--- convenience methods + bool subscribe(string channel) {return setSubscribe(channel);} + bool unsubscribe(string channel) {return setUnsubscribe(channel);} + + //--- only for ZMQ_XSUB + SOCKOPT_BOOL(XpubVerbose,ZMQ_XPUB_VERBOSE) + SOCKOPT_BOOL(XpubVerboser,ZMQ_XPUB_VERBOSER) + SOCKOPT_BOOL(XpubManual,ZMQ_XPUB_MANUAL) + SOCKOPT_BOOL(XpubNoDrop,ZMQ_XPUB_NODROP) // also for ZMQ_PUB + SOCKOPT_SET_BYTES(XpubWelcomeMessage,ZMQ_XPUB_WELCOME_MSG) + + SOCKOPT_BOOL(InvertMatching,ZMQ_INVERT_MATCHING) //--- only for ZMQ_PUB, ZMQ_XPUB, ZMQ_SUB + + //--- only for ZMQ_ROUTER + SOCKOPT_SET_BOOL(RouterHandover,ZMQ_ROUTER_HANDOVER) + SOCKOPT_SET_BOOL(RouterMandatory,ZMQ_ROUTER_MANDATORY) + SOCKOPT_SET_BOOL(RouterRaw,ZMQ_ROUTER_RAW) + + //--- only for ZMQ_STREAM + SOCKOPT_SET_BOOL(StreamNotify,ZMQ_STREAM_NOTIFY) + + //--- only for ZMQ_ROUTER, ZMQ_STREAM + SOCKOPT_SET_BYTES(ConnectRid,ZMQ_CONNECT_RID) + + //--- only for ZMQ_REP, ZMQ_REQ, ZMQ_ROUTER, ZMQ_DEALER + SOCKOPT_BYTES(Identity,ZMQ_IDENTITY,255) + + //--- tcp + SOCKOPT(int,TcpKeepAlive,ZMQ_TCP_KEEPALIVE) + SOCKOPT(int,TcpKeepAliveCount,ZMQ_TCP_KEEPALIVE_CNT) + SOCKOPT(int,TcpKeepAliveIdle,ZMQ_TCP_KEEPALIVE_IDLE) + SOCKOPT(int,TcpKeepAliveInterval,ZMQ_TCP_KEEPALIVE_INTVL) + SOCKOPT(int,TcpMaxRetransmitTimeout,ZMQ_TCP_MAXRT) + + SOCKOPT(int,TypeOfService,ZMQ_TOS) // IP_TOS + + //--- ZMQ_TCP_ACCEPT_FILTER + //--- ZMQ_IPC_FILTER_GID + //--- ZMQ_IPC_FILTER_PID + //--- ZMQ_IPC_FILTER_UID + //--- are deprecated in favor of ZAP API and ip address whitelisting/blacklisting + SOCKOPT_NTSTR(ZapDomain,ZMQ_ZAP_DOMAIN) + + //--- only for vmci transport + SOCKOPT(ulong,VmciBufferSize,ZMQ_VMCI_BUFFER_SIZE) // bytes + SOCKOPT(ulong,VmciBufferMinSize,ZMQ_VMCI_BUFFER_MIN_SIZE) // bytes + SOCKOPT(ulong,VmciBufferMaxSize,ZMQ_VMCI_BUFFER_MAX_SIZE) // bytes + SOCKOPT(int,VmciConnectTimeout,ZMQ_VMCI_CONNECT_TIMEOUT) // milliseconds + }; +//+------------------------------------------------------------------+ +//| The option value is a string with predefined byte length | +//| | +//| If it is a NULL-terminated string without predefined length, | +//| The situation is tricky: we do not know the length of the option | +//| value beforehand, but the function does not return the correct | +//| one, either. So the only option is to guess. | +//| | +//| Here we adopt the solution of the Java binding. We just guess | +//| that the length of a NULL-terminated string option is less than | +//| 1024. So hopefully, it is the case. | +//+------------------------------------------------------------------+ +bool SocketOptions::getStringOption(int option,string &value,size_t length) + { + char buf[]; + ArrayResize(buf,(int)length); + bool res=getOption(option,buf,length); + if(res) + { + value=StringFromUtf8(buf); + } + ArrayFree(buf); + return res; + } +//+------------------------------------------------------------------+ +//| The ending means that the converted buffer contains the ending | +//| null. | +//+------------------------------------------------------------------+ +bool SocketOptions::setStringOption(int option,const string value,bool ending) + { + char buf[]; + StringToUtf8(value,buf,ending); + int len = ArraySize(buf); + bool res=setOption(option,buf,len); + ArrayFree(buf); + return res; + } +//+------------------------------------------------------------------+ diff --git a/Include/Zmq/Z85.mqh b/Include/Zmq/Z85.mqh new file mode 100644 index 0000000..15dc980 --- /dev/null +++ b/Include/Zmq/Z85.mqh @@ -0,0 +1,140 @@ +//+------------------------------------------------------------------+ +//| Z85.mqh | +//| Copyright 2016, Li Ding | +//| dingmaotu@hotmail.com | +//+------------------------------------------------------------------+ +#property copyright "Copyright 2016, Li Ding" +#property link "dingmaotu@hotmail.com" +#property strict + +#include "Common.mqh" + +#import "libzmq.dll" +// Encode data with Z85 encoding. Returns 0(NULL) if failed +intptr_t zmq_z85_encode(char &str[],const uchar &data[],size_t size); + +// Decode data with Z85 encoding. Returns 0(NULL) if failed +intptr_t zmq_z85_decode(uchar &dest[],const char &str[]); + +// Generate z85-encoded public and private keypair with tweetnacl/libsodium +int zmq_curve_keypair(char &z85_public_key[],char &z85_secret_key[]); + +// Derive the z85-encoded public key from the z85-encoded secret key +int zmq_curve_public(char &z85_public_key[],const char &z85_secret_key[]); +#import +//+------------------------------------------------------------------+ +//| Z85 encoding/decoding | +//+------------------------------------------------------------------+ +class Z85 + { +public: + static bool encode(string &secret,const uchar &data[]); + static bool decode(const string secret,uchar &data[]); + + static string encode(string data); + static string decode(string secret); + + static bool generateKeyPair(uchar &publicKey[],uchar &secretKey[]); + static bool derivePublic(uchar &publicKey[],const uchar &secretKey[]); + + static bool generateKeyPair(string &publicKey,string &secretKey); + static string derivePublic(const string secretKey); + }; +//+------------------------------------------------------------------+ +//| data must have size multiple of 4 | +//+------------------------------------------------------------------+ +bool Z85::encode(string &secret,const uchar &data[]) + { + int size=ArraySize(data); + if(size%4 != 0) return false; + + char str[]; + ArrayResize(str,(int)(1.25*size+1)); + + intptr_t res=zmq_z85_encode(str,data,size); + if(res == 0) return false; + secret = StringFromUtf8(str); + return true; + } +//+------------------------------------------------------------------+ +//| secret must be multiples of 5 | +//+------------------------------------------------------------------+ +bool Z85::decode(const string secret,uchar &data[]) + { + int len=StringLen(secret); + if(len%5 != 0) return false; + + char str[]; + StringToUtf8(secret,str); + ArrayResize(data,(int)(0.8*len)); + return 0 != zmq_z85_decode(data,str); + } +//+------------------------------------------------------------------+ +//| data length should be multiples of 4 and only ascii is supported | +//+------------------------------------------------------------------+ +string Z85::encode(string data) + { + char str[]; + StringToUtf8(data,str,false); + string res; + if(encode(res,str)) + return res; + else + return ""; + } +//+------------------------------------------------------------------+ +//| secret must be multiples of 5 | +//+------------------------------------------------------------------+ +string Z85::decode(string secret) + { + uchar data[]; + decode(secret,data); + return StringFromUtf8(data); + } +//+------------------------------------------------------------------+ +//| | +//+------------------------------------------------------------------+ +bool Z85::generateKeyPair(uchar &publicKey[],uchar &secretKey[]) + { + ArrayResize(publicKey,41); + ArrayResize(secretKey,41); + return 0==zmq_curve_keypair(publicKey, secretKey); + } +//+------------------------------------------------------------------+ +//| | +//+------------------------------------------------------------------+ +bool Z85::derivePublic(uchar &publicKey[],const uchar &secretKey[]) + { + ArrayResize(publicKey,41); + return 0==zmq_curve_public(publicKey, secretKey); + } +//+------------------------------------------------------------------+ +//| | +//+------------------------------------------------------------------+ +bool Z85::generateKeyPair(string &publicKey,string &secretKey) + { + uchar sec[],pub[]; + bool res=generateKeyPair(pub,sec); + if(res) + { + secretKey=StringFromUtf8(sec); + publicKey=StringFromUtf8(pub); + } + ArrayFree(sec); + ArrayFree(pub); + return res; + } +//+------------------------------------------------------------------+ +//| | +//+------------------------------------------------------------------+ +string Z85::derivePublic(const string secrect) + { + uchar sec[],pub[]; + StringToUtf8(secrect,sec); + derivePublic(pub,sec); + string pubstr=StringFromUtf8(pub); + ArrayFree(sec); + ArrayFree(pub); + return pubstr; + } +//+------------------------------------------------------------------+ diff --git a/Include/Zmq/Zmq.mqh b/Include/Zmq/Zmq.mqh new file mode 100644 index 0000000..332139e Binary files /dev/null and b/Include/Zmq/Zmq.mqh differ diff --git a/Include/Zmq/ZmqMsg.mqh b/Include/Zmq/ZmqMsg.mqh new file mode 100644 index 0000000..cd63a97 --- /dev/null +++ b/Include/Zmq/ZmqMsg.mqh @@ -0,0 +1,120 @@ +//+------------------------------------------------------------------+ +//| ZmqMsg.mqh | +//| Copyright 2016, Li Ding | +//| dingmaotu@hotmail.com | +//+------------------------------------------------------------------+ +#property copyright "Copyright 2016, Li Ding" +#property link "dingmaotu@hotmail.com" +#property strict +#include "Common.mqh" +//+------------------------------------------------------------------+ +//| 0MQ Message struct | +//+------------------------------------------------------------------+ +// align sizeof(intptr_t) +// = 8 on 64bit +// = 4 on 32bit +// hopefully MetaQuotes do malloc structs +// aligned on pointer address boundaries +struct zmq_msg_t + { + uchar _[64]; + }; + +#import "libzmq.dll" +int zmq_msg_init(zmq_msg_t &msg); +int zmq_msg_init_size(zmq_msg_t &msg,size_t size); +// As mt4 can not provide a zmq_free_fn, and can not let +// zmq library own the array data, copying is always needed. +// Therefore this function will not be used by this binding. +// int zmq_msg_init_data(zmq_msg_t &msg,uchar &data[], +// int size,int ffn,int hint); +int zmq_msg_close(zmq_msg_t &msg); +int zmq_msg_move(zmq_msg_t &dest,zmq_msg_t &src); +int zmq_msg_copy(zmq_msg_t &dest,zmq_msg_t &src); +// char * +intptr_t zmq_msg_data(zmq_msg_t &msg); +int zmq_msg_size(zmq_msg_t &msg); +int zmq_msg_more(zmq_msg_t &msg); +int zmq_msg_get(zmq_msg_t &msg,int property); +int zmq_msg_set(zmq_msg_t &msg,int property,int optval); +// const char * +intptr_t zmq_msg_gets(zmq_msg_t &msg,const char &property[]); +#import +//+------------------------------------------------------------------+ +//| Wraps a zmq_msg_t | +//+------------------------------------------------------------------+ +struct ZmqMsg: public zmq_msg_t + { +protected: + int get(int property) {return zmq_msg_get(this,property);} + bool set(int property,int value) {return 0==zmq_msg_set(this,property,value);} + intptr_t data() {return zmq_msg_data(this);} +public: + ZmqMsg() {zmq_msg_init(this);} + ZmqMsg(int size) {if(0!=zmq_msg_init_size(this,size)){Debug("Failed to init size msg: insufficient space");}} + ZmqMsg(string data,bool nullterminated=false); + ~ZmqMsg() {if(0!=zmq_msg_close(this)){Debug("Failed to close msg");}} + + size_t size() {return zmq_msg_size(this);} + + void getData(uchar &data[]); + string getData(); + void setData(const uchar &data[]); + + bool more() {return 1==zmq_msg_more(this);} + + bool copy(ZmqMsg &msg) {return 0 == zmq_msg_copy(this, msg);} + bool move(ZmqMsg &msg) {return 0 == zmq_msg_move(this, msg);} + + string meta(const string property); + }; +//+------------------------------------------------------------------+ +//| Initialize a utf-8 string message | +//+------------------------------------------------------------------+ +ZmqMsg::ZmqMsg(string data,bool nullterminated) + { + uchar array[]; + StringToUtf8(data,array,nullterminated); + int size=ArraySize(array); + zmq_msg_init_size(this,size); + setData(array); + } +//+------------------------------------------------------------------+ +//| Get message data as bytes array | +//+------------------------------------------------------------------+ +void ZmqMsg::getData(uchar &data[]) + { + size_t size=size(); + intptr_t src=data(); + ArrayResize(data,(int)size); + ArrayFromPointer(data,src); + } +//+------------------------------------------------------------------+ +//| Get message data as utf-8 string | +//+------------------------------------------------------------------+ +string ZmqMsg::getData() + { + size_t size=size(); + intptr_t psz=data(); + return StringFromUtf8Pointer(psz,(int)size); + } +//+------------------------------------------------------------------+ +//| copy data to message internal storage | +//+------------------------------------------------------------------+ +void ZmqMsg::setData(const uchar &data[]) + { + intptr_t dest=data(); + ArrayToPointer(data,dest); + } +//+------------------------------------------------------------------+ +//| Wraps zmq_msg_gets: get metadata associated with the msg | +//+------------------------------------------------------------------+ +string ZmqMsg::meta(const string property) + { + uchar buf[]; + StringToUtf8(property,buf); + intptr_t ref=zmq_msg_gets(this,buf); + ArrayFree(buf); + return StringFromUtf8Pointer(ref); + } +//+------------------------------------------------------------------+ diff --git a/README.md b/README.md index a424bbc..a6d5957 100644 --- a/README.md +++ b/README.md @@ -1,2 +1,41 @@ # mql-zmq + ZMQ binding for the MQL language (both 32bit MT4 and 64bit MT5) + +## Introduction + +This is a complete binding of the [ZeroMQ](http://zeromq.org/) library +for the MQL4/5 language provided by MetaTrader4/5. + +Traders with programming abilities have always wanted a messaging solution +like ZeroMQ, simple and powerful, far better than the PIPE trick as +suggested by the official articles. However, bindings for MQL were either outdated or not complete (a lot of functionalities are missing). This binding is based on latest 4.2 version of the library, and provided all functionalities as specified in the API documentation. + +This binding tries to remain compatible between MQL4/5. Users of both versions can use this binding, with a single set of headers. MQL4 and MQL5 are basically the same in that they are merged in recent versions. The difference is in the runtime environment (MetaTrader5 is 64bit by default, while MetaTrader4 is 32bit). The trading system is also different, but it is no concern of this binding. + +## Files + +This binding contains three sets of files: + +1. The binding itself is in the `Include/Zmq` directory. + +2. The testing scripts and zmq guide examples are in `Scripts` directory. The script files are mq4 by default, but you can change the extension to mq5 to use them in MetaTrader5. + +3. Precompiled DLLs of both 64bit (`Library/MT5`) and 32bit (`Library/MT4`) ZeroMQ and libsodium are provided. Copy the corresponding DLLs to the `Library` folder of your MetaTrader terminal. If you are using MT5 32bit, use the 32bit version from `Library/MT4`. The DLLs require that you have the latest Visual C++ runtime (2015). *Note* that these DLLs are compiled from official sources, without any modification. You can compile your own if you don't trust these binaries. The `libsodium.dll` is copied from the official binary release. If you want to support security mechanisms other than `curve`, or you want to use transports like OpenPGM, you need to compile your own DLL. + +## About string encoding + +MQL strings are Win32 UNICODE strings (basically 2-byte UTF-16). In this binding all strings are converted to utf-8 strings before sending to the dll layer. The ZmqMsg supports a constructor from MQL strings, the default is _NOT_ null-terminated. + +## Usage + +You can find a simple test script in `Scripts/Test`, and you can find examples of the official guide in Scripts/ZeroMQGuideExamples. I intend to translate all examples to this binding, but now only the hello world example is provided. I will gradually add those examples. Of course forking this binding if you are interested and welcome to send pull requests. + +## TODO + +1. Write more tests. +2. Add more examples from the official ZMQ guide. + +## Changes + +2016-12-27: Released 1.0. \ No newline at end of file diff --git a/Scripts/Test/TestZmq.mq4 b/Scripts/Test/TestZmq.mq4 new file mode 100644 index 0000000..25e5d64 --- /dev/null +++ b/Scripts/Test/TestZmq.mq4 @@ -0,0 +1,93 @@ +//+------------------------------------------------------------------+ +//| TestZmq.mq4 | +//| Copyright 2016, Li Ding | +//| dingmaotu@hotmail.com | +//+------------------------------------------------------------------+ +#property copyright "Copyright 2016, Li Ding" +#property link "dingmaotu@hotmail.com" +#property version "1.00" +#property strict + +#include +//+------------------------------------------------------------------+ +//| Script program start function | +//+------------------------------------------------------------------+ +void OnStart() + { +//--- Test capabilities + Print(">>> Testing capabilities"); + Print("0) Zmq version is [",Zmq::getVersion(),"]"); + Print("1) Supports ipc:// protocol: [",Zmq::hasIpc(),"]"); + Print("2) Supports pgm:// protocol: [",Zmq::hasPgm(),"]"); + Print("3) Supports norm:// protocol: [",Zmq::hasNorm(),"]"); + Print("4) Supports tipc:// protocol: [",Zmq::hasTipc(),"]"); + Print("5) Supports curve security: [",Zmq::hasCurve(),"]"); + Print("6) Supports gssapi security: [",Zmq::hasGssApi(),"]"); + Print(">>> End testing capabilities"); + +//--- Test Z85 encoding/decoding + Print(">>> Testing Z85 encoding/decoding"); + + string data="12345678"; + Print("1) Original data is: ",data); + + string secret=Z85::encode(data); + Print("2) Encrypted value is: ",secret); + + string decoded=Z85::decode(secret); + Print("3) Decoded: ",decoded); + Print("4) Decoded value is equal to original: [",decoded==data,"]"); + Print(">>> End testing Z85 encoding/decoding"); + +//--- Test atomic counters + Print(">>> Testing atomic counters"); + AtomicCounter counter; + Print("1) Initial value should be 0: [",counter.get()==0,"]"); + counter.set(5); + Print("2) Counter set to 5: [",counter.get()==5,"]"); + counter.increase(); + Print("3) Increased value should be 6: [",counter.get()==6,"]"); + counter.decrease(); + Print("4) Decreased value should be 5: [",counter.get()==5,"]"); + Print(">>> End testing atomic counters"); + +//--- Test context + Print(">>> Testing context"); + Context context; + Print("1) Default IO threads should be ZMQ_IO_THREADS_DFLT: [",context.getIoThreads()==ZMQ_IO_THREADS_DFLT,"]"); + Print(">>> trying to set io threads to 2"); + context.setIoThreads(2); + Print("2) Now IO threads should be 2: [",context.getIoThreads()==2,"]"); + Print("3) Socket limit: [",context.getSocketLimit(),"]"); + Print("4) Max sockets: [",context.getMaxSockets(),"]"); + Print("5) Max message size: [",context.getMaxMessageSize(),"]"); + Print(">>> End testing context"); + + Socket s(context,ZMQ_REP); + string addr="inproc://abc"; + if(!s.bind(addr)) + { + Debug(StringFormat("Error binding %s: %s",addr,Zmq::errorMessage())); + } + else + { + Debug(StringFormat("Success binded %s",addr)); + } + string endpoint; + s.getLastEndpoint(endpoint); + Print("Last endpoint is: [",endpoint,"]"); + + string principal; + s.getGssApiPrincipal(principal); + Print("Principal is [",principal,"]"); + +//--- Test curve + Print(">>> Testing curve"); + string genpub,gensec; + Z85::generateKeyPair(genpub,gensec); + Print("1) Generated public key: [",genpub,"]"); + Print("1) Generated private key: [",gensec,"]"); + Print("2) Derive public key from secret key: [",Z85::derivePublic(gensec)==genpub,"]"); + Print(">>> End testing curve"); + } +//+------------------------------------------------------------------+ diff --git a/Scripts/ZeroMQGuideExamples/Chapter1/HelloWorldClient.mq4 b/Scripts/ZeroMQGuideExamples/Chapter1/HelloWorldClient.mq4 new file mode 100644 index 0000000..83fc118 Binary files /dev/null and b/Scripts/ZeroMQGuideExamples/Chapter1/HelloWorldClient.mq4 differ diff --git a/Scripts/ZeroMQGuideExamples/Chapter1/HelloWorldServer.mq4 b/Scripts/ZeroMQGuideExamples/Chapter1/HelloWorldServer.mq4 new file mode 100644 index 0000000..4906578 --- /dev/null +++ b/Scripts/ZeroMQGuideExamples/Chapter1/HelloWorldServer.mq4 @@ -0,0 +1,42 @@ +//+------------------------------------------------------------------+ +//| HelloWorldServer.mq5 | +//| Copyright 2016, Li Ding | +//| dingmaotu@hotmail.com | +//+------------------------------------------------------------------+ +#property copyright "Copyright 2016, Li Ding" +#property link "dingmaotu@hotmail.com" +#property version "1.00" + +#include +//+------------------------------------------------------------------+ +//| Hello World server in MQL | +//| Binds REP socket to tcp://*:5555 | +//| Expects "Hello" from client, replies with "World" | +//+------------------------------------------------------------------+ +void OnStart() + { + Context context; + Socket socket(context,ZMQ_REP); + + socket.bind("tcp://*:5555"); + + while(true) + { + ZmqMsg request; + + // Wait for next request from client + + // MetaTrader note: this will block the script thread + // and if you try to terminate this script, MetaTrader + // will hang (and crash if you force closing it) + socket.recv(request); + Print("Receive Hello"); + + Sleep(1000); + + ZmqMsg reply("World"); + // Send reply back to client + socket.send(reply); + } + } +//+------------------------------------------------------------------+