1 Commits

Author SHA1 Message Date
Ding Li 01afabc408 Refactored poll support. Add Chapter 2 examples from the official ZMQ guide. 2017-07-18 16:40:22 +08:00
7 changed files with 309 additions and 24 deletions
+32 -16
View File
@@ -9,14 +9,6 @@
#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
@@ -64,7 +56,17 @@ struct zmq_pollitem_t
#define ZMQ_POLLPRI 8
#define ZMQ_POLLITEMS_DFLT 16
//--- fd is SOCKET on Win32, which is defined as UINT_PTR
struct PollItem
{
intptr_t socket;
uintptr_t fd;
short events;
short revents;
bool hasInput() const {return(revents&ZMQ_POLLIN)!=0;}
bool hasOutput() const {return(revents&ZMQ_POLLOUT)!=0;}
};
#import "libzmq.dll"
//+------------------------------------------------------------------+
//| Sockets |
@@ -87,7 +89,7 @@ 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);
int zmq_poll(PollItem &items[],int nitems,long timeout);
//+------------------------------------------------------------------+
//| Message proxying |
//+------------------------------------------------------------------+
@@ -123,8 +125,8 @@ public:
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);
void register(PollItem &pollitem,bool read=false,bool write=false);
void register(PollItem &pollitems[],int index,bool read=false,bool write=false);
//--- monitor socket events
bool monitor(string addr,int events);
@@ -134,7 +136,10 @@ public:
static bool proxySteerable(Socket *frontend,Socket *backend,Socket *capture,Socket *control);
//--- poll
static int poll(zmq_pollitem_t &arr[],long timeout);
static int poll(PollItem &arr[],long timeout);
//--- fill a poll item for this socket
void fillPollItem(PollItem &item,short events);
};
//+------------------------------------------------------------------+
//| |
@@ -242,7 +247,7 @@ bool Socket::monitor(string addr,int events)
//+------------------------------------------------------------------+
//| |
//+------------------------------------------------------------------+
void Socket::register(zmq_pollitem_t &pollitem,bool read=false,bool write=false)
void Socket::register(PollItem &pollitem,bool read=false,bool write=false)
{
ZeroMemory(pollitem);
pollitem.socket=m_ref;
@@ -252,7 +257,7 @@ void Socket::register(zmq_pollitem_t &pollitem,bool read=false,bool write=false)
//+------------------------------------------------------------------+
//| |
//+------------------------------------------------------------------+
void Socket::register(zmq_pollitem_t &pollitems[],int index,bool read=false,bool write=false)
void Socket::register(PollItem &pollitems[],int index,bool read=false,bool write=false)
{
ZeroMemory(pollitems[index]);
pollitems[index].socket=m_ref;
@@ -281,10 +286,21 @@ bool Socket::proxySteerable(Socket *frontend,Socket *backend,Socket *capture,Soc
return 0==zmq_proxy_steerable(frontend_ref, backend_ref, capture_ref, control_ref);
}
//+------------------------------------------------------------------+
//| |
//| poll for events. timeout is milliseconds (-1 for indefinite wait |
//| and 0 for immediate return) |
//+------------------------------------------------------------------+
int Socket::poll(zmq_pollitem_t &arr[],long timeout)
int Socket::poll(PollItem &arr[],long timeout)
{
return zmq_poll(arr,ArraySize(arr),timeout);
}
//+------------------------------------------------------------------+
//| fill a poll item for this socket |
//+------------------------------------------------------------------+
void Socket::fillPollItem(PollItem &item,short events)
{
item.socket=m_ref;
item.fd=0;
item.events=events;
item.revents=0;
}
//+------------------------------------------------------------------+
+36 -8
View File
@@ -7,11 +7,19 @@ ZMQ binding for the MQL language (both 32bit MT4 and 64bit MT5)
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 (mostly toy projects and only basic features are implemented). This binding is based on latest 4.2 version of the library, and provides all functionalities as specified in the API documentation.
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 (mostly toy projects and only basic features are implemented). This
binding is based on latest 4.2 version of the library, and provides 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.
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
@@ -19,13 +27,27 @@ 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.
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.
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.
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.
## Notes on context creation
@@ -46,7 +68,11 @@ globally, and in a manner not easily recognized by humans, for example:
## 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.
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.
Here is a sample from `HelloWorldServer.mq4`:
@@ -92,6 +118,8 @@ void OnStart()
## Changes
* 2017-07-18: Released 1.3: Refactored poll support. Add Chapter 2 examples from
the official ZMQ guide.
* 2017-06-08: Released 1.2: Fix GlobalHandle bug. Add rebuild method to ZmqMsg.
Complete all examples in ZMQ Guide Chanpter 1.
* 2017-05-26: Released 1.1: add the ability to share a ZMQ context globally in a terminal
@@ -0,0 +1,54 @@
//+------------------------------------------------------------------+
//| MSPoller.mq4 |
//| Copyright 2017, Li Ding |
//| dingmaotu@126.com |
//+------------------------------------------------------------------+
#property copyright "Copyright 2017, Li Ding"
#property link "dingmaotu@126.com"
#property version "1.00"
#property strict
#include <Zmq/Zmq.mqh>
//+------------------------------------------------------------------+
//| Reading from multiple sockets in MQL (adapted from C++ version) |
//| This version uses zmq_poll() |
//| |
//| Olivier Chamoux <olivier.chamoux@fr.thalesgroup.com> |
//+------------------------------------------------------------------+
void OnStart()
{
//---
Context context;
// Connect to task ventilator
Socket receiver(context,ZMQ_PULL);
receiver.connect("tcp://localhost:5557");
// Connect to weather server
Socket subscriber(context,ZMQ_SUB);
subscriber.connect("tcp://localhost:5556");
subscriber.subscribe("10001 ");
// Initialize poll set
PollItem items[2];
receiver.fillPollItem(items[0],ZMQ_POLLIN);
subscriber.fillPollItem(items[1],ZMQ_POLLIN);
// Process messages from both sockets
while(!IsStopped())
{
ZmqMsg message;
//--- MQL Note: To handle Script exit properly, we set a timeout of 500 ms instead of infinite wait
Socket::poll(items,500);
if(items[0].hasInput())
{
receiver.recv(message);
// Process task
}
if(items[1].hasInput())
{
subscriber.recv(message);
// Process weather update
}
}
}
//+------------------------------------------------------------------+
@@ -0,0 +1,63 @@
//+------------------------------------------------------------------+
//| RRBroker.mq4 |
//| Copyright 2017, Li Ding |
//| dingmaotu@126.com |
//+------------------------------------------------------------------+
#property copyright "Copyright 2017, Li Ding"
#property link "dingmaotu@126.com"
#property version "1.00"
#property strict
#include <Zmq/Zmq.mqh>
//+------------------------------------------------------------------+
//| Simple request-reply broker in MQL (adapted from C++ version) |
//| |
//| Olivier Chamoux <olivier.chamoux@fr.thalesgroup.com> |
//+------------------------------------------------------------------+
void OnStart()
{
// Prepare our context and sockets
Context context;
Socket frontend(context,ZMQ_ROUTER);
Socket backend(context,ZMQ_DEALER);
frontend.bind("tcp://*:5559");
backend.bind("tcp://*:5560");
// Initialize poll set
PollItem items[2];
frontend.fillPollItem(items[0],ZMQ_POLLIN);
backend.fillPollItem(items[1],ZMQ_POLLIN);
// Switch messages between sockets
while(!IsStopped())
{
ZmqMsg message;
bool more=false; // Multipart detection
Socket::poll(items,500);
if(items[0].hasInput())
{
// Process all parts of the message
do
{
frontend.recv(message);
more=message.more();
backend.send(message,false,more);
}
while(more);
}
if(items[1].hasInput())
{
// Process all parts of the message
do
{
backend.recv(message);
more=message.more();
frontend.send(message,false,more);
}
while(more);
}
}
}
//+------------------------------------------------------------------+
@@ -0,0 +1,36 @@
//+------------------------------------------------------------------+
//| RRClient.mq4 |
//| Copyright 2017, Li Ding |
//| dingmaotu@126.com |
//+------------------------------------------------------------------+
#property copyright "Copyright 2017, Li Ding"
#property link "dingmaotu@126.com"
#property version "1.00"
#property strict
#include <Zmq/Zmq.mqh>
//+------------------------------------------------------------------+
//| Request-reply client in MQL (adapted from C++ version) |
//| Connects REQ socket to tcp://localhost:5559 |
//| Sends "Hello" to server, expects "World" back |
//| |
//| Olivier Chamoux <olivier.chamoux@fr.thalesgroup.com> |
//+------------------------------------------------------------------+
void OnStart()
{
Context context;
Socket requester(context,ZMQ_REQ);
requester.connect("tcp://localhost:5559");
for(int request=0; request<10; request++)
{
ZmqMsg message("Hello");
requester.send(message);
ZmqMsg reply;
requester.recv(reply,true);
Print("Received reply ",reply.getData());
}
}
//+------------------------------------------------------------------+
@@ -0,0 +1,41 @@
//+------------------------------------------------------------------+
//| RRWorker.mq4 |
//| Copyright 2017, Li Ding |
//| dingmaotu@126.com |
//+------------------------------------------------------------------+
#property copyright "Copyright 2017, Li Ding"
#property link "dingmaotu@126.com"
#property version "1.00"
#property strict
#include <Zmq/Zmq.mqh>
//+------------------------------------------------------------------+
//| Request-reply service in MQL (adapted from C++ version) |
//| Connects REP socket to tcp://localhost:5560 |
//| Expects "Hello" from client, replies with "World" |
//| |
//| Olivier Chamoux <olivier.chamoux@fr.thalesgroup.com> |
//+------------------------------------------------------------------+
void OnStart()
{
Context context;
Socket responder(context,ZMQ_REP);
responder.connect("tcp://localhost:5560");
while(!IsStopped())
{
// Wait for next request from client
ZmqMsg req;
responder.recv(req);
Print("Received request: ",req.getData());
// Do some 'work'
Sleep(1000);
ZmqMsg reply("World");
// Send reply back to client
responder.send(reply);
}
}
//+------------------------------------------------------------------+
@@ -0,0 +1,47 @@
//+------------------------------------------------------------------+
//| WUProxy.mq4 |
//| Copyright 2017, Li Ding |
//| dingmaotu@126.com |
//+------------------------------------------------------------------+
#property copyright "Copyright 2017, Li Ding"
#property link "dingmaotu@126.com"
#property version "1.00"
#property strict
#include <Zmq/Zmq.mqh>
//+------------------------------------------------------------------+
//| Weather proxy device MQL (adapted from C++) |
//| |
//| Olivier Chamoux <olivier.chamoux@fr.thalesgroup.com> |
//+------------------------------------------------------------------+
void OnStart()
{
//---
Context context;
// This is where the weather server sits
Socket frontend(context,ZMQ_XSUB);
frontend.connect("tcp://192.168.55.210:5556");
// This is our public endpoint for subscribers
Socket backend(context,ZMQ_XPUB);
backend.bind("tcp://10.1.1.0:8100");
// Subscribe on everything
frontend.subscribe("");
// Shunt messages out to our own subscribers
while(!IsStopped())
{
// Process all parts of the message
ZmqMsg message;
bool more;
do
{
frontend.recv(message);
more=message.more();
backend.send(message,false,more);
}
while(more);
}
}
//+------------------------------------------------------------------+