preview
Beyond REST and ZeroMQ: Building a gRPC/Protocol Buffers Bridge for Real-Time MetaTrader 5–Python Inference

Beyond REST and ZeroMQ: Building a gRPC/Protocol Buffers Bridge for Real-Time MetaTrader 5–Python Inference

MetaTrader 5 — Integration |
78 0
Olamide Daniel Adebayo
Olamide Daniel Adebayo

Introduction

Bridging a live MetaTrader 5 EA to a Python model server usually goes through the same progression: file-polling (race conditions on Windows file locks), then WebRequest/REST (HTTP overhead, JSON parsing, no schema), then ZeroMQ (fixes the transport but people still hand-roll JSON strings over it, so the schema problem never actually goes away). Neither REST nor ZeroMQ solves the real issue: nothing enforces that the feature vector MQL5 sends matches what the Python model expects, and a silent field mismatch trades on garbage instead of failing loudly.

This article builds the missing contract: a Protocol Buffers schema for the MetaTrader 5–Python boundary. The transport behaves like gRPC, but MQL5 cannot speak gRPC natively (no HTTP/2 client, no stream multiplexing, and no protoc code generation). So instead of overclaiming "gRPC in MQL5," this builds a length-prefixed Protobuf-over-TCP framing that a small Python shim relays into a genuine grpc.aio service. The Python side gets real gRPC — streaming, health checks, the actual protocol. The MQL5 side gets a strongly-typed, schema-enforced message format with none of the ambiguity of string-concatenated JSON.

Seven core files ship with this article (plus benchmark and cache-generation scripts), each walked through with its actual code below — what it implements, the key mechanism, and why that approach over the alternatives.

Scope note: this is not "gRPC in MQL5" in the strict sense — MQL5 cannot open an HTTP/2 stream. What you get is a Protobuf-encoded, length-prefixed TCP channel that a Python shim bridges into a real gRPC backend. This distinction is only restated once more, in the Architecture section, and not repeated after that.


Architecture and Schema

The system has three layers:

  • the MQL5 client (CGrpcInferenceClient, built on CProtoWriter/CProtoReader), which opens a raw TCP socket to the shim, sends a length-prefixed Protobuf message, and blocks for a framed response with timeout and retry;
  • the Python shim (grpc_shim_server.py), an asyncio TCP server that deserializes those frames into real Protobuf objects and forwards them over an actual grpc.aio channel to the model server;
  • the Python model server (model_server.py), a genuine grpc.aio.server() implementing the schema below, unary and streaming.

Why not skip the shim and call the model server directly from MQL5? That would mean hand-implementing HTTP/2 framing and HPACK in MQL5 — a large undertaking with real correctness risk — or calling something "gRPC" that isn't. The shim is a small, auditable piece of Python that keeps the rest of the system honest about what it actually is.

MetaTrader 5–Python Inference Bridge: Three-Layer Architecture

Fig. 1. The three-layer bridge: MQL5 client with a hand-rolled Protobuf encoder, a Python TCP-to-gRPC shim, and a genuine grpc.aio model server. Only the shim-to-model-server hop is real HTTP/2 gRPC.

The schema, inference.proto, is the single source of truth every field number in the MQL5 client mirrors by hand, since there's no MQL5 protoc plugin:

syntax = "proto3";

package mt5inference;

// FeatureVector: the input payload built by GrpcInferenceEA.mq5's
// BuildFeatureVector(). "values" is a packed repeated float, encoded
// on the wire as one length-delimited blob of concatenated 4-byte
// little-endian floats rather than one tag per value (see
// CProtoWriter::WritePackedFloatArray in ProtobufWire.mqh).
message FeatureVector {
  string symbol = 1;
  string timeframe = 2;
  int64 timestamp = 3;
  repeated float values = 4;
}

// InferenceRequest: what CGrpcInferenceClient::Predict() sends per bar.
// "features" is an embedded FeatureVector message (nested/length-
// delimited), and "model_version" lets the server reject a request
// built against a stale contract before it ever reaches the model.
message InferenceRequest {
  string request_id = 1;
  FeatureVector features = 2;
  string model_version = 3;
}

// InferenceResponse: decoded by CProtoReader inside
// CGrpcInferenceClient::ParseResponse(). A non-empty "error" is this
// project's success/failure flag -- see SInferenceResult.valid in
// GrpcInferenceClient.mqh.
message InferenceResponse {
  string request_id = 1;
  float probability = 2;
  int32 class_label = 3;
  float latency_ms = 4;
  string model_version = 5;
  string error = 6;
}

// InferenceService: implemented for real by model_server.py using
// grpc.aio. grpc_shim_server.py is a client of this service on MT5's
// behalf -- MQL5 itself never calls these RPCs directly.
service InferenceService {
  rpc Predict(InferenceRequest) returns (InferenceResponse);
  rpc PredictStream(stream InferenceRequest) returns (stream InferenceResponse);
}

FeatureVector.values is a packed repeated float — proto3 concatenates it into one length-delimited blob rather than tagging each element, which matters for an 8-float vector sent every bar. InferenceRequest carries the embedded FeatureVector plus a model version the server can reject if stale. InferenceResponse's error field doubles as the success/failure flag. InferenceService declares both unary Predict and streaming PredictStream; the MQL5 client only exercises the unary path today, but the streaming RPC is already live on the Python side.


MQL5 Protobuf Implementation: ProtobufWire.mqh

MQL5 has no protoc plugin, so ProtobufWire.mqh hand-implements the wire format: a tag (field number shifted left 3 bits, OR'd with a 3-bit wire type) followed by a payload, with integers encoded as base-128 varints — 7 value bits per byte, top bit as a continuation flag, so small field numbers and lengths cost a single byte. CProtoWriter builds outgoing messages; CProtoReader parses incoming ones. Field numbers matter more than they look: they're literally what gets written onto the wire in place of a name, so every constant in the client that mirrors inference.proto exists purely to keep the hand-written encoder in sync with the schema.

Writer — buffer and core primitive. The internal byte buffer grows geometrically (double plus slack) rather than to the exact size needed, so appending one byte at a time stays close to amortized O(1) instead of triggering a full copy on every single append. WriteVarint is the primitive every other method funnels through; WriteTag just packs the field number and wire type into one varint call.

The buffer and its append primitive:

//+------------------------------------------------------------------+
//|                                             ProtobufWire.mqh     |
//|                                       https://www.mql5.com       |
//+------------------------------------------------------------------+
#property strict

//+------------------------------------------------------------------+
//| Protobuf wire-type constants                                     |
//+------------------------------------------------------------------+
#define PB_WIRETYPE_VARINT   0     // int32/int64/bool/enum on the wire
#define PB_WIRETYPE_FIXED64  1     // double, fixed64, sfixed64
#define PB_WIRETYPE_LENDELIM 2     // string, bytes, embedded message, packed repeated
#define PB_WIRETYPE_FIXED32  5     // float, fixed32, sfixed32

//+------------------------------------------------------------------+
//| CProtoWriter                                                     |
//| Minimal hand-rolled Protobuf writer. MQL5 has no protoc plugin,  |
//| so every method below encodes exactly what protoc would generate |
//| for the matching .proto field type -- documented per method.     |
//+------------------------------------------------------------------+
class CProtoWriter
  {
private:
   //+------------------------------------------------------------------+
   //| Internal growable byte buffer                                    |
   //+------------------------------------------------------------------+
   uchar             m_buf[];         // raw output bytes accumulated so far
   int               m_len;           // number of bytes actually used in m_buf

   //+------------------------------------------------------------------+
   //| EnsureCapacity                                                   |
   //+------------------------------------------------------------------+
   void              EnsureCapacity(const int extra)
     {
      //--- grow geometrically (double + slack) instead of by exact need,
      //--- so repeated single-byte appends don't cause O(n^2) resizing
      if(m_len+extra>ArraySize(m_buf))
         ArrayResize(m_buf,MathMax(m_len+extra,ArraySize(m_buf)*2+64));
     }

   //+------------------------------------------------------------------+
   //| AppendByte                                                       |
   //+------------------------------------------------------------------+
   void              AppendByte(const uchar b)
     {
      EnsureCapacity(1);
      m_buf[m_len++]=b;
     }

Construction and byte extraction:

public:
   //+------------------------------------------------------------------+
   //| CProtoWriter (constructor)                                       |
   //+------------------------------------------------------------------+
                     CProtoWriter(void) : m_len(0) { ArrayResize(m_buf,64); }

   //+------------------------------------------------------------------+
   //| Reset                                                            |
   //+------------------------------------------------------------------+
   void              Reset(void) { m_len=0; } // reuse the buffer across calls, no realloc

   //+------------------------------------------------------------------+
   //| Size                                                             |
   //+------------------------------------------------------------------+
   int               Size(void) const { return m_len; }

   //+------------------------------------------------------------------+
   //| GetBytes                                                         |
   //| Copies exactly m_len bytes out -- m_buf itself may be larger     |
   //| than m_len because of the geometric growth in EnsureCapacity.    |
   //+------------------------------------------------------------------+
   void              GetBytes(uchar &out[])
     {
      ArrayResize(out,m_len);
      for(int i=0;i<m_len;i++)
         out[i]=m_buf[i];
     }

The varint and tag encoders that everything else calls:

   //+------------------------------------------------------------------+
   //| WriteVarint                                                      |
   //| Unsigned LEB128 varint: 7 payload bits per byte, MSB=1 means     |
   //| "more bytes follow". Every other Write*Field method below funnels|
   //| through this for its length or numeric payload.                  |
   //+------------------------------------------------------------------+
   void              WriteVarint(ulong value)
     {
      while(value>=0x80)
        {
         //--- emit low 7 bits with continuation bit set, then shift
         AppendByte((uchar)((value & 0x7F) | 0x80));
         value>>=7;
        }
      AppendByte((uchar)value); // final byte, continuation bit is 0
     }

   //+------------------------------------------------------------------+
   //| WriteTag                                                         |
   //| Builds the (field_number<<3 | wire_type) tag varint that         |
   //| precedes every field on the wire, per the Protobuf spec.         |
   //+------------------------------------------------------------------+
   void              WriteTag(const int field_number,const int wire_type)
     {
      WriteVarint((ulong)((field_number<<3) | wire_type));
     }

Writer — scalar and composite fields. Int32 and int64 both encode as varints — Protobuf doesn't distinguish them on the wire, only in how the decoded value gets typed afterward. Floats use a small hand-rolled union type (UFloatBytes) declared near the top of the file — MQL5 doesn't support C-style anonymous inline unions declared inside a function body the way C/C++ does, so a named union type is required, the same way a struct would be. This writes fixed32 values as little-endian bytes and assumes a little-endian host (x86/x64). Re-check this if the code is ever ported. Strings are length-delimited UTF-8; note the -1 that trims StringToCharArray's trailing null terminator, which Protobuf's length prefix must not count — skipping that subtraction silently makes every string field one byte too long.

The feature vector itself goes through WritePackedFloatArray, proto3's default packed encoding for repeated scalar fields: one tag and one length-prefixed blob of concatenated 4-byte floats for the whole 8-element array, rather than a separate tag per element. Embedded messages (InferenceRequest.features) are simpler than they sound — the submessage is fully serialized first with its own writer instance, and the resulting bytes are then written as one more length-delimited field, exactly how Predict() builds the outer request around the feature vector later in this article.

Int64 and int32 fields (both encode as the same varint):

   //+------------------------------------------------------------------+
   //| WriteInt64Field                                                  |
   //+------------------------------------------------------------------+
   void              WriteInt64Field(const int field_number,const long value)
     {
      WriteTag(field_number,PB_WIRETYPE_VARINT);
      WriteVarint((ulong)value);
     }

   //+------------------------------------------------------------------+
   //| WriteInt32Field                                                  |
   //| int32 has no distinct wire representation -- it is promoted to   |
   //| int64 and encoded the same way, matching protoc's own behavior.  |
   //+------------------------------------------------------------------+
   void              WriteInt32Field(const int field_number,const int value)
     {
      WriteInt64Field(field_number,(long)value);
     }

The named union type the float methods below rely on:

//+------------------------------------------------------------------+
//| UFloatBytes                                                      |
//+------------------------------------------------------------------+
union UFloatBytes
  {
   float             f;      // the value as a 32-bit IEEE-754 float
   uchar             b[4];   // the same 4 bytes, addressable individually
  };

Float encoding, using that union:

   //+------------------------------------------------------------------+
   //| FloatToBytesLE                                                   |
   //| Protobuf fixes 32-bit numeric fields as little-endian on the     |
   //| wire. The union re-interprets the float's 4 raw bytes without any|
   //| numeric conversion -- this assumes a little-endian host CPU,     |
   //| true for the x86/x64 terminals this EA targets.                  |
   //+------------------------------------------------------------------+
   static void       FloatToBytesLE(const float value,uchar &out[])
     {
      ArrayResize(out,4);
      UFloatBytes u;
      u.f=value;
      out[0]=u.b[0]; out[1]=u.b[1]; out[2]=u.b[2]; out[3]=u.b[3];
     }

   //+------------------------------------------------------------------+
   //| WriteFloatField                                                  |
   //+------------------------------------------------------------------+
   void              WriteFloatField(const int field_number,const float value)
     {
      WriteTag(field_number,PB_WIRETYPE_FIXED32);
      uchar b[4];
      FloatToBytesLE(value,b);
      for(int i=0;i<4;i++)
         AppendByte(b[i]);
     }

String fields (length-delimited UTF-8):

   //+------------------------------------------------------------------+
   //| WriteStringField                                                 |
   //| Strings are length-delimited: tag, then a varint byte count,     |
   //| then the raw UTF-8 bytes. StringToCharArray appends a trailing   |
   //| null terminator that Protobuf does not want on the wire, so the  |
   //| count is taken one short of the returned array size.             |
   //+------------------------------------------------------------------+
   void              WriteStringField(const int field_number,const string value)
     {
      uchar sb[];
      int n=StringToCharArray(value,sb,0,-1,CP_UTF8)-1; // drop trailing '\0'
      if(n<0)
         n=0;
      WriteTag(field_number,PB_WIRETYPE_LENDELIM);
      WriteVarint((ulong)n);
      for(int i=0;i<n;i++)
         AppendByte(sb[i]);
     }

The packed float array used for the feature vector:

   //+------------------------------------------------------------------+
   //| WritePackedFloatArray                                            |
   //| proto3 packs repeated scalar-numeric fields by default: ONE      |
   //| length-delimited blob of concatenated fixed32 floats, rather     |
   //| than a separate tag per array element. This is the encoding for  |
   //| FeatureVector.values in inference.proto.                         |
   //+------------------------------------------------------------------+
   void              WritePackedFloatArray(const int field_number,const float &values[])
     {
      int n=ArraySize(values);
      WriteTag(field_number,PB_WIRETYPE_LENDELIM);
      WriteVarint((ulong)(n*4)); // byte length of the packed blob, not element count
      for(int i=0;i<n;i++)
        {
         uchar b[4];
         FloatToBytesLE(values[i],b);
         for(int j=0;j<4;j++)
            AppendByte(b[j]);
        }
     }

Embedding one message inside another:

   //+------------------------------------------------------------------+
   //| WriteEmbeddedMessage                                             |
   //| Nested messages (e.g. InferenceRequest.features) are encoded by  |
   //| serializing the submessage first with its own writer instance,   |
   //| then treating the resulting bytes as one length-delimited blob.  |
   //| The caller is responsible for serializing submsg beforehand.     |
   //+------------------------------------------------------------------+
   void              WriteEmbeddedMessage(const int field_number,const uchar &submsg[])
     {
      int n=ArraySize(submsg);
      WriteTag(field_number,PB_WIRETYPE_LENDELIM);
      WriteVarint((ulong)n);
      for(int i=0;i<n;i++)
         AppendByte(submsg[i]);
     }
  };

Reader — the forward-compatibility mechanism. CProtoReader mirrors the writer end to end: Attach copies the byte array (not a reference, so the reader stays valid even if the caller's array is resized mid-parse) and resets a read cursor; ReadVarint and ReadTag decode the same bit-packing in reverse, with a shift-overflow guard capping at 63 bits so a corrupted or truncated frame can't spin the parser indefinitely waiting for a continuation bit that never clears. The one method worth calling out specifically is SkipField: given only a wire type, it advances past a field's payload without knowing or caring what the field means. That's what makes the parser forward-compatible — if model_server.py is ever updated to add a 7th response field, an EA compiled against this article's code keeps working, silently ignoring the new field instead of failing to parse the message at all. ParseResponse's default case, covered in the next section, relies on exactly this.

The reader's constructor and cursor state:

//+------------------------------------------------------------------+
//| CProtoReader                                                     |
//| Minimal hand-rolled Protobuf reader used to parse the            |
//| InferenceResponse bytes that come back from the Python shim.     |
//+------------------------------------------------------------------+
class CProtoReader
  {
private:
   uchar             m_buf[];        // local copy of the bytes being parsed
   int               m_len;          // total length of m_buf in use
   int               m_pos;          // current read cursor

public:
   //+------------------------------------------------------------------+
   //| Attach                                                           |
   //| Copies the caller's byte array in and resets the read cursor.    |
   //| A copy (not a reference) is kept deliberately so the reader      |
   //| stays valid even if the caller's array is resized afterward.     |
   //+------------------------------------------------------------------+
   void              Attach(const uchar &data[],const int len)
     {
      ArrayResize(m_buf,len);
      for(int i=0;i<len;i++)
         m_buf[i]=data[i];
      m_len=len;
      m_pos=0;
     }

End-of-buffer check and the varint decoder:

   //+------------------------------------------------------------------+
   //| Eof                                                              |
   //+------------------------------------------------------------------+
   bool              Eof(void) const { return m_pos>=m_len; }

   //+------------------------------------------------------------------+
   //| ReadVarint                                                       |
   //| Mirror of CProtoWriter::WriteVarint. The shift>63 guard exists   |
   //| because a corrupted or truncated frame could otherwise never set |
   //| the continuation bit to 0, spinning through the whole buffer.    |
   //+------------------------------------------------------------------+
   bool              ReadVarint(ulong &value)
     {
      value=0;
      int shift=0;
      while(m_pos<m_len)
        {
         uchar b=m_buf[m_pos++];
         value|=((ulong)(b & 0x7F))<<shift;
         if((b & 0x80)==0)
            return true;      // continuation bit clear -- varint complete
         shift+=7;
         if(shift>63)
            return false;     // malformed-varint guard, see Edge Cases section
        }
      return false;           // ran out of bytes mid-varint
     }

Tag parsing and the unknown-field skip:

   //+------------------------------------------------------------------+
   //| ReadTag                                                          |
   //+------------------------------------------------------------------+
   bool              ReadTag(int &field_number,int &wire_type)
     {
      ulong tag;
      if(!ReadVarint(tag))
         return false;
      field_number=(int)(tag>>3);
      wire_type=(int)(tag & 0x07);
      return true;
     }

   //+------------------------------------------------------------------+
   //| SkipField                                                        |
   //| Skips a field's payload based purely on its wire type, without   |
   //| knowing what the field means. This is what makes the parser      |
   //| forward-compatible: an unknown field number from a newer server  |
   //| is skipped cleanly instead of breaking the whole parse.          |
   //+------------------------------------------------------------------+
   bool              SkipField(const int wire_type)
     {
      ulong tmp;
      switch(wire_type)
        {
         case PB_WIRETYPE_VARINT:
            return ReadVarint(tmp);
         case PB_WIRETYPE_FIXED64:
            m_pos+=8;
            return(m_pos<=m_len);
         case PB_WIRETYPE_FIXED32:
            m_pos+=4;
            return(m_pos<=m_len);
         case PB_WIRETYPE_LENDELIM:
           {
            ulong n;
            if(!ReadVarint(n))
               return false;
            m_pos+=(int)n;
            return(m_pos<=m_len);
           }
         default:
            return false;     // unrecognized wire type -- cannot safely skip
        }
     }

Fixed-width floats and length-delimited reads:

   //+------------------------------------------------------------------+
   //| ReadFloat                                                        |
   //| Inverse of CProtoWriter::FloatToBytesLE -- reinterprets 4 raw    |
   //| little-endian bytes back into a float via the same union trick.  |
   //+------------------------------------------------------------------+
   bool              ReadFloat(float &value)
     {
      if(m_pos+4>m_len)
         return false;
      UFloatBytes u;
      u.b[0]=m_buf[m_pos]; u.b[1]=m_buf[m_pos+1]; u.b[2]=m_buf[m_pos+2]; u.b[3]=m_buf[m_pos+3];
      m_pos+=4;
      value=u.f;
      return true;
     }

   //+------------------------------------------------------------------+
   //| ReadLengthDelimited                                              |
   //+------------------------------------------------------------------+
   bool              ReadLengthDelimited(uchar &out[])
     {
      ulong n;
      if(!ReadVarint(n))
         return false;
      if(m_pos+(int)n>m_len)
         return false;        // declared length would overrun the buffer
      ArrayResize(out,(int)n);
      for(int i=0;i<(int)n;i++)
         out[i]=m_buf[m_pos+i];
      m_pos+=(int)n;
      return true;
     }

String and int32 convenience wrappers:

   //+------------------------------------------------------------------+
   //| ReadString                                                       |
   //+------------------------------------------------------------------+
   bool              ReadString(string &value)
     {
      uchar raw[];
      if(!ReadLengthDelimited(raw))
         return false;
      value=CharArrayToString(raw,0,ArraySize(raw),CP_UTF8);
      return true;
     }

   //+------------------------------------------------------------------+
   //| ReadInt32                                                        |
   //+------------------------------------------------------------------+
   bool              ReadInt32(int &value)
     {
      ulong v;
      if(!ReadVarint(v))
         return false;
      value=(int)v;
      return true;
     }


MQL5 Client and EA Integration

Here's the actual byte layout the client sends, since it's the reference worth having on hand while reading this section:

Length-Prefixed Protobuf Framing on the Wire

Fig. 2. The byte layout MetaTrader 5 actually sends: a 4-byte length prefix, then the serialized message with per-field tags, and a zoomed view of how a packed repeated-float field is laid out inside it.

GrpcInferenceClient.mqh turns the encoder into a network call, in four parts. Constants and result struct: field-number #define constants that must mirror inference.proto exactly, since a mismatch here silently mistags a field on the wire.

//+------------------------------------------------------------------+
//|                                      GrpcInferenceClient.mqh     |
//|                                       https://www.mql5.com       |
//+------------------------------------------------------------------+
#property strict
#include "ProtobufWire.mqh"

//+------------------------------------------------------------------+
//| InferenceRequest field numbers -- MUST mirror inference.proto    |
//+------------------------------------------------------------------+
#define GRPC_FIELD_REQ_ID        1     // request_id       (string)
#define GRPC_FIELD_FEATURES      2     // features         (FeatureVector, embedded)
#define GRPC_FIELD_MODEL_VERSION 3     // model_version    (string)

//+------------------------------------------------------------------+
//| FeatureVector field numbers -- MUST mirror inference.proto       |
//+------------------------------------------------------------------+
#define GRPC_FV_FIELD_SYMBOL     1     // symbol           (string)
#define GRPC_FV_FIELD_TIMEFRAME  2     // timeframe        (string)
#define GRPC_FV_FIELD_TIMESTAMP  3     // timestamp        (int64)
#define GRPC_FV_FIELD_VALUES     4     // values           (repeated float, packed)

//+------------------------------------------------------------------+
//| InferenceResponse field numbers -- MUST mirror inference.proto   |
//+------------------------------------------------------------------+
#define GRPC_RESP_FIELD_REQ_ID       1     // request_id     (string)
#define GRPC_RESP_FIELD_PROBABILITY  2     // probability    (float)
#define GRPC_RESP_FIELD_CLASS_LABEL  3     // class_label    (int32)
#define GRPC_RESP_FIELD_LATENCY_MS   4     // latency_ms     (float)
#define GRPC_RESP_FIELD_MODEL_VER    5     // model_version  (string)
#define GRPC_RESP_FIELD_ERROR        6     // error          (string)

//+------------------------------------------------------------------+
//| SInferenceResult                                                 |
//| Decoded, ready-to-use view of an InferenceResponse. valid is set |
//| by ParseResponse() based on whether error came back empty.       |
//+------------------------------------------------------------------+
struct SInferenceResult
  {
   string            request_id;      // echoed back from the request
   float             probability;     // model output in [0,1]
   int               class_label;     // 1 = long signal, 0 = short signal, -1 = unset
   float             latency_ms;      // server-side inference time, for diagnostics
   string            model_version;   // model version the server actually used
   string            error;           // non-empty means the call failed logically
   bool              valid;           // convenience flag mirroring (error=="")
  };

Length-prefixed framing: a 4-byte big-endian length header precedes every payload in both directions. SendFramed writes header then payload, checking both writes actually sent the full byte count. RecvExact loops because SocketRead can return fewer bytes than requested on a single call; RecvFramed reads the header first and rejects any declared length over 16 MB before trusting it for the payload read — this bound protects against a corrupted length field causing a runaway allocation or an indefinite wait.

private:
   //+------------------------------------------------------------------+
   //| Connection state and configuration                               |
   //+------------------------------------------------------------------+
   int               m_socket;        // MQL5 socket handle, or INVALID_HANDLE
   string            m_host;          // shim host, e.g. "127.0.0.1"
   int               m_port;          // shim port, e.g. 50151
   int               m_timeout_ms;    // per-call receive deadline
   int               m_max_retries;   // attempts beyond the first before giving up

   //+------------------------------------------------------------------+
   //| SendFramed                                                       |
   //| Writes the 4-byte big-endian length prefix, then the payload     |
   //| itself. Both SocketSend calls are checked against the byte count |
   //| they were asked to send -- a short write means the transport is  |
   //| in a bad state and the caller should reconnect, not keep going.  |
   //+------------------------------------------------------------------+
   bool              SendFramed(const uchar &payload[])
     {
      int n=ArraySize(payload);
      uchar header[4];
      header[0]=(uchar)((n>>24)&0xFF);
      header[1]=(uchar)((n>>16)&0xFF);
      header[2]=(uchar)((n>>8)&0xFF);
      header[3]=(uchar)(n&0xFF);

      if(SocketSend(m_socket,header,4)!=4)
         return false;
      if(n>0 && SocketSend(m_socket,payload,n)!=n)
         return false;
      return true;
     }

   //+------------------------------------------------------------------+
   //| RecvExact                                                        |
   //| Blocks until exactly n bytes have arrived or deadline_ms has     |
   //| elapsed. SocketRead can return fewer bytes than requested on a   |
   //| single call, so this loop keeps accumulating into out[] rather   |
   //| than assuming one call returns the whole frame. Unlike a         |
   //| self-growing stream API, SocketRead expects its destination      |
   //| buffer pre-sized to the amount being requested -- chunk[] is     |
   //| resized to exactly how many bytes are both available and still   |
   //| needed before each call.                                         |
   //+------------------------------------------------------------------+
   bool              RecvExact(uchar &out[],const int n,const int deadline_ms)
     {
      ArrayResize(out,n);
      int received=0;
      uint start=GetTickCount();
      while(received<n)
        {
         if((int)(GetTickCount()-start)>deadline_ms)
            return false;                       // deadline exceeded
         uint avail=SocketIsReadable(m_socket);
         if(avail==0)
           {
            Sleep(1);
            continue;                            // nothing to read yet
           }
         uint want=(uint)(n-received);
         if(avail<want)
            want=avail;                          // don't ask for more than is buffered
         uchar chunk[];
         ArrayResize(chunk,(int)want);
         int got=SocketRead(m_socket,chunk,want,10);
         if(got<=0)
           {
            Sleep(1);
            continue;
           }
         for(int i=0;i<got;i++)
            out[received+i]=chunk[i];
         received+=got;
        }
      return true;
     }

   //+------------------------------------------------------------------+
   //| RecvFramed                                                       |
   //| Reads the 4-byte length header first, sanity-checks it against   |
   //| a hard upper bound (see Edge Cases section), then reads exactly  |
   //| that many payload bytes.                                         |
   //+------------------------------------------------------------------+
   bool              RecvFramed(uchar &payload[],const int deadline_ms)
     {
      uchar header[];
      if(!RecvExact(header,4,deadline_ms))
         return false;
      int n=(header[0]<<24)|(header[1]<<16)|(header[2]<<8)|header[3];
      if(n<0 || n>16*1024*1024)
         return false;                           // corrupt-frame sanity guard
      return RecvExact(payload,n,deadline_ms);
     }

Response parsing with unknown-field skipping: ParseResponse walks the response fields via a switch, dispatching known fields to typed setters and falling through to SkipField for anything else — the forward-compatibility path described in the previous section. result.valid is simply "was error empty," a single source of truth for success.

   //+------------------------------------------------------------------+
   //| ParseResponse                                                    |
   //| Walks the InferenceResponse bytes field by field. The default    |
   //| case in the switch calls SkipField() rather than failing, so a   |
   //| response with fields this client doesn't recognize still parses  |
   //| the fields it does know about.                                   |
   //+------------------------------------------------------------------+
   bool              ParseResponse(const uchar &data[],SInferenceResult &result)
     {
      CProtoReader r;
      r.Attach(data,ArraySize(data));

      result.request_id="";
      result.probability=0.0f;
      result.class_label=-1;
      result.latency_ms=0.0f;
      result.model_version="";
      result.error="";

      while(!r.Eof())
        {
         int fn,wt;
         if(!r.ReadTag(fn,wt))
            break;

         switch(fn)
           {
            case GRPC_RESP_FIELD_REQ_ID:      r.ReadString(result.request_id);    break;
            case GRPC_RESP_FIELD_PROBABILITY: r.ReadFloat(result.probability);    break;
            case GRPC_RESP_FIELD_CLASS_LABEL: r.ReadInt32(result.class_label);    break;
            case GRPC_RESP_FIELD_LATENCY_MS:  r.ReadFloat(result.latency_ms);     break;
            case GRPC_RESP_FIELD_MODEL_VER:   r.ReadString(result.model_version); break;
            case GRPC_RESP_FIELD_ERROR:       r.ReadString(result.error);         break;
            default:                          r.SkipField(wt);                   break;
           }
        }

      result.valid=(StringLen(result.error)==0);
      return result.valid;
     }

Predict flow — reconnect, retry, backoff: Predict builds the embedded FeatureVector, wraps it into the outer InferenceRequest, and then runs a send/receive loop with exponential backoff (50ms × 2^attempt, so a transient shim restart isn't hammered by a tight retry loop while a one-off glitch still resolves quickly). On any failure — connect, send, or receive — it calls Disconnect() before retrying, because the socket's state is unknown after a failed exchange. If the retry budget is exhausted, it reports a transport error through result.error — the same field a logical (feature-mismatch) failure uses, so OnTick() checks both failure modes identically, with no special-casing needed.

   //+------------------------------------------------------------------+
   //| Predict                                                          |
   //| Full request/response cycle: serialize the embedded FeatureVector|
   //| first, wrap it into the outer InferenceRequest, then send/receive|
   //| with retry. Each retry uses exponential backoff (50ms * 2^n) so a|
   //| transient shim restart doesn't get hammered by a tight loop.     |
   //+------------------------------------------------------------------+
   bool              Predict(const string request_id,
                              const string symbol,
                              const string timeframe,
                              const long timestamp,
                              const float &features[],
                              const string model_version,
                              SInferenceResult &result)
     {
      result.valid=false;

      //--- build FeatureVector submessage first, then embed it below
      CProtoWriter fvw;
      fvw.WriteStringField(GRPC_FV_FIELD_SYMBOL,symbol);
      fvw.WriteStringField(GRPC_FV_FIELD_TIMEFRAME,timeframe);
      fvw.WriteInt64Field(GRPC_FV_FIELD_TIMESTAMP,timestamp);
      fvw.WritePackedFloatArray(GRPC_FV_FIELD_VALUES,features);
      uchar fv_bytes[];
      fvw.GetBytes(fv_bytes);

      //--- build the outer InferenceRequest around the embedded bytes
      CProtoWriter reqw;
      reqw.WriteStringField(GRPC_FIELD_REQ_ID,request_id);
      reqw.WriteEmbeddedMessage(GRPC_FIELD_FEATURES,fv_bytes);
      reqw.WriteStringField(GRPC_FIELD_MODEL_VERSION,model_version);
      uchar req_bytes[];
      reqw.GetBytes(req_bytes);

      int attempt=0;
      while(attempt<=m_max_retries)
        {
         if(!IsConnected() && !Connect())
           {
            attempt++;
            Sleep(50*(1<<attempt));           // exponential backoff before retry
            continue;
           }

         if(!SendFramed(req_bytes))
           {
            Disconnect();                     // socket is in an unknown state, drop it
            attempt++;
            Sleep(50*(1<<attempt));
            continue;
           }

         uchar resp_bytes[];
         if(!RecvFramed(resp_bytes,m_timeout_ms))
           {
            Disconnect();
            attempt++;
            Sleep(50*(1<<attempt));
            continue;
           }

         return ParseResponse(resp_bytes,result);
        }

      result.error="transport failure after retries";
      return false;
     }

GrpcInferenceEA.mq5 drives all of this. OnInit() cross-checks that BuildFeatureVector()'s actual output width matches the NUM_FEATURES constant before anything else runs, aborting loudly on a mismatch rather than sending a malformed vector; it also resolves once whether this run is inside the Strategy Tester (MQLInfoInteger(MQL_TESTER)), which decides whether the EA connects a live socket or loads an offline-recorded cache — covered in the Testing section.

OnInit():

//+------------------------------------------------------------------+
//| OnInit                                                           |
//| Runs the feature-contract cross-check before doing anything else:|
//| BuildFeatureVector()'s actual output width must equal the        |
//| NUM_FEATURES constant, or initialization aborts loudly rather    |
//| than silently sending a mismatched vector to the model server.   |
//+------------------------------------------------------------------+
int OnInit()
  {
   float probe[];
   int built_width=BuildFeatureVector(probe);
   if(built_width!=NUM_FEATURES)
     {
      PrintFormat("[FEATURE_CONTRACT_MISMATCH] NUM_FEATURES=%d but BuildFeatureVector() produced %d columns. Aborting init.",
                  NUM_FEATURES,built_width);
      return(INIT_PARAMETERS_INCORRECT);
     }

   ExtAtrHandle=iATR(_Symbol,_Period,InpAtrPeriod);
   if(ExtAtrHandle==INVALID_HANDLE)
     {
      Print("[INIT_FAIL] Could not create ATR handle.");
      return(INIT_FAILED);
     }

   ExtClient.Configure(InpServerHost,InpServerPort,InpSocketTimeoutMs,InpMaxRetries);
   if(!ExtClient.Connect())
      //--- non-fatal: the retry/backoff logic inside Predict() will
      //--- attempt to reconnect on the first real tick regardless
      Print("[GRPC_WARN] Initial connect to inference shim failed; will retry on first tick.");

   ExtTrade.SetExpertMagicNumber(990211);
   return(INIT_SUCCEEDED);
  }

OnDeinit():

//+------------------------------------------------------------------+
//| OnDeinit                                                         |
//+------------------------------------------------------------------+
void OnDeinit(const int reason)
  {
   ExtClient.Disconnect();
   if(ExtAtrHandle!=INVALID_HANDLE)
      IndicatorRelease(ExtAtrHandle);
  }

BuildFeatureVector() computes 8 values every bar: three lagged one-bar returns, candle range/body/wick ratios, and normalized ATR — every ratio divided by the current close so features stay comparable regardless of the instrument's absolute price level (a 5-point range means something very different on a $50 stock than on XAUUSD near $4,300). The same function runs from both OnInit() and OnTick(), which matters: it guarantees the startup width check tests the exact code path that produces the real feature vector at runtime, not a hand-maintained duplicate that could silently drift out of sync.

//+------------------------------------------------------------------+
//| BuildFeatureVector                                               |
//| Computes the 8-column feature row sent to the model: three       |
//| lagged returns, candle range/body/wick ratios, and normalized    |
//| ATR. All ratios are normalized by the current close so the model |
//| sees scale-invariant inputs regardless of instrument price level.|
//| Also called from OnInit() with an empty probe array purely to    |
//| report the column count for the feature-contract check -- in that|
//| case there may not be enough warmed-up bars yet, so the guards   |
//| below return a correctly-sized zero vector rather than failing.  |
//+------------------------------------------------------------------+
int BuildFeatureVector(float &out[])
  {
   ArrayResize(out,NUM_FEATURES);

   if(Bars(_Symbol,_Period)<InpAtrPeriod+5)
     {
      ArrayInitialize(out,0.0f);
      return NUM_FEATURES;              // not enough history yet, e.g. at OnInit
     }

   double atr[];
   ArraySetAsSeries(atr,true);
   if(ExtAtrHandle==INVALID_HANDLE || CopyBuffer(ExtAtrHandle,0,0,1,atr)<=0)
     {
      ArrayInitialize(out,0.0f);
      return NUM_FEATURES;
     }

   MqlRates rates[];
   ArraySetAsSeries(rates,true);
   if(CopyRates(_Symbol,_Period,0,6,rates)<6)
     {
      ArrayInitialize(out,0.0f);
      return NUM_FEATURES;
     }

   //--- three lagged one-bar returns, most recent first
   double ret1 = (rates[0].close-rates[1].close)/rates[1].close;
   double ret2 = (rates[1].close-rates[2].close)/rates[2].close;
   double ret3 = (rates[2].close-rates[3].close)/rates[3].close;

   //--- candle-shape ratios, all normalized by the current close
   double range = (rates[0].high-rates[0].low)/rates[0].close;
   double body  = MathAbs(rates[0].close-rates[0].open)/rates[0].close;
   double atr_norm    = atr[0]/rates[0].close;
   double upper_wick  = (rates[0].high-MathMax(rates[0].open,rates[0].close))/rates[0].close;
   double lower_wick  = (MathMin(rates[0].open,rates[0].close)-rates[0].low)/rates[0].close;

   out[0]=(float)ret1;
   out[1]=(float)ret2;
   out[2]=(float)ret3;
   out[3]=(float)range;
   out[4]=(float)body;
   out[5]=(float)atr_norm;
   out[6]=(float)upper_wick;
   out[7]=(float)lower_wick;

   return NUM_FEATURES;
  }

OnTick() is rate-limited to once per new bar — every call past that point is a real network round trip, and calling on every tick (which can fire dozens of times per second on a fast-moving pair) would both waste the latency budget covered in the Results section and make the retry logic above trigger far more often under load. It builds the feature vector, re-checks its width defensively, and calls either Predict() or the cache lookup depending on ExtUsingCache, resolved once at OnInit(). ManageSignal() gates a trade on probability >= InpProbThreshold and sizes stop-loss/take-profit off the live ATR rather than a fixed pip distance, since XAUUSD's volatility varies too much across sessions for a fixed stop to make sense — a stop that's reasonable during a quiet Asian session can be far too tight during a US-session news spike.

OnTick():

//+------------------------------------------------------------------+
//| OnTick                                                           |
//| Rate-limited to once per new M5 bar -- every call below is a real|
//| network round trip, so calling on every tick would both waste the|
//| latency budget and make the client's retry/backoff far more      |
//| likely to trigger under load. Runtime feature width is checked   |
//| again here (not just at OnInit) in case history conditions change|
//| the vector shape mid-run.                                        |
//+------------------------------------------------------------------+
void OnTick()
  {
   datetime cur_bar_time=iTime(_Symbol,_Period,0);
   if(cur_bar_time==ExtLastBarTime)
      return;                             // still the same bar, nothing to do
   ExtLastBarTime=cur_bar_time;

   float features[];
   int width=BuildFeatureVector(features);
   if(width!=NUM_FEATURES)
     {
      PrintFormat("[FEATURE_CONTRACT_MISMATCH] Runtime width %d != NUM_FEATURES %d. Skipping tick.",width,NUM_FEATURES);
      return;
     }

   string req_id=StringFormat("%s-%I64d",_Symbol,(long)TimeCurrent());
   SInferenceResult result;

   bool ok=ExtClient.Predict(req_id,_Symbol,EnumToString(_Period),(long)TimeCurrent(),
                              features,InpModelVersion,result);

   if(!ok || !result.valid)
     {
      PrintFormat("[GRPC_ERROR] Inference call failed: %s",result.error);
      return;
     }

   if(result.model_version!=InpModelVersion)
      //--- non-fatal by design: the server is authoritative on which
      //--- model actually ran, this is a visibility log line only
      PrintFormat("[MODEL_VERSION_DRIFT] Server responded with model_version=%s, expected %s.",
                  result.model_version,InpModelVersion);

   ManageSignal(result);
  }

ManageSignal():

//+------------------------------------------------------------------+
//| ManageSignal                                                     |
//| Converts a validated InferenceResult into a threshold-gated,     |
//| ATR-sized market order. Single-position gating keeps the example |
//| easy to follow -- production use would likely add per-symbol     |
//| exposure limits and a cooldown after a losing trade.             |
//+------------------------------------------------------------------+
void ManageSignal(const SInferenceResult &result)
  {
   if(PositionSelect(_Symbol))
      return;                             // already in a position, skip this signal

   double atr[];
   ArraySetAsSeries(atr,true);
   if(CopyBuffer(ExtAtrHandle,0,0,1,atr)<=0)
      return;
   double atr_val=atr[0];
   if(atr_val<=0)
      return;                             // guard against a zero/invalid ATR read

   double price_ask=SymbolInfoDouble(_Symbol,SYMBOL_ASK);
   double price_bid=SymbolInfoDouble(_Symbol,SYMBOL_BID);

   if(result.class_label==1 && result.probability>=InpProbThreshold)
     {
      double sl=price_ask-InpAtrMultiplier*atr_val;
      double tp=price_ask+InpAtrMultiplier*atr_val;
      ExtTrade.Buy(InpLotSize,_Symbol,price_ask,sl,tp,"grpc-inference-long");
     }
   else if(result.class_label==0 && result.probability>=InpProbThreshold)
     {
      double sl=price_bid+InpAtrMultiplier*atr_val;
      double tp=price_bid-InpAtrMultiplier*atr_val;
      ExtTrade.Sell(InpLotSize,_Symbol,price_bid,sl,tp,"grpc-inference-short");
     }
  }


Python Relay and Model Server

grpc_shim_server.py is the piece that makes "gRPC bridge" honest rather than a rebrand of a plain socket relay. It reads MetaTrader 5's length-prefixed frames with read_frame — a Python-side mirror of the MQL5 client's own framing, using struct.unpack(">I", header) to fix the byte order as big-endian, matching the manual byte-shifting done by hand on the MQL5 side — deserializes them into real inference_pb2.InferenceRequest objects, and forwards them over a genuine grpc.aio channel to the model server. This is the one hop in the whole system that's unambiguously real gRPC: HTTP/2 framing, handled entirely by the grpc library rather than anything hand-rolled here.

A grpc.aio.AioRpcError (model server unreachable, timed out, or erroring) is caught and turned into an ordinary error-populated response, so the MQL5 side only ever has to check result.error once, regardless of which layer actually failed. One channel is created at startup and reused across every MetaTrader 5 connection the shim ever accepts, since grpc.aio.insecure_channel is explicitly designed to be safe for concurrent use and expensive to establish — this is also what lets several Strategy Tester agents share the shim during a multi-agent optimization run with no extra pooling logic needed here.

"""
grpc_shim_server.py

Bridges MT5's length-prefixed Protobuf-over-TCP framing (MQL5 has no HTTP/2
stack, so it cannot speak real gRPC) to a genuine grpc.aio channel talking to
model_server.py. This file is intentionally small (~60 lines of logic) and
does exactly one job: translate framing, nothing more -- it never inspects
or modifies message contents, only relays the raw serialized bytes.

Wire framing expected from MT5 (see CGrpcInferenceClient.mqh):
    [4 bytes big-endian length][serialized InferenceRequest protobuf bytes]
Response is framed identically with a serialized InferenceResponse.
"""
import asyncio
import logging
import struct

import grpc

import inference_pb2
import inference_pb2_grpc

SHIM_HOST = "127.0.0.1"
SHIM_PORT = 50151            # matches InpServerPort default in GrpcInferenceEA.mq5
MODEL_SERVER_TARGET = "127.0.0.1:50152"

logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s")
log = logging.getLogger("grpc_shim")


async def read_frame(reader: asyncio.StreamReader) -> bytes:
    """Reads exactly one framed message: a 4-byte big-endian length prefix
    followed by that many payload bytes. readexactly() blocks until the
    full amount is available or raises IncompleteReadError on disconnect --
    this is the Python-side mirror of CGrpcInferenceClient::RecvExact() on
    the MQL5 side, and both sides agree on big-endian for this length field
    specifically (independent of the little-endian floats inside the
    Protobuf payload itself)."""
    header = await reader.readexactly(4)
    (length,) = struct.unpack(">I", header)
    if length > 16 * 1024 * 1024:
        # same 16 MB sanity bound enforced on the MQL5 side in RecvFramed();
        # keeping both bounds identical avoids one side accepting a frame
        # the other would have rejected.
        raise ValueError(f"refusing oversized frame: {length} bytes")
    return await reader.readexactly(length)


def write_frame(writer: asyncio.StreamWriter, payload: bytes) -> None:
    """Writes the mirror-image framing: big-endian length, then payload,
    in a single writer.write() call so the two pieces can never be
    interleaved with another coroutine's write on the same connection."""
    writer.write(struct.pack(">I", len(payload)) + payload)


async def handle_client(reader: asyncio.StreamReader, writer: asyncio.StreamWriter,
                         stub: inference_pb2_grpc.InferenceServiceStub) -> None:
    """One coroutine per connected MT5 client -- asyncio.start_server spins
    this up automatically per accepted socket, which is what lets multiple
    Strategy Tester agents connect to this shim concurrently during a
    multi-agent optimization run without any extra pooling logic here."""
    peer = writer.get_extra_info("peername")
    log.info("MT5 client connected: %s", peer)
    try:
        while True:
            try:
                raw = await read_frame(reader)
            except asyncio.IncompleteReadError:
                break  # client closed the socket

            # deserialize into a REAL Protobuf message object (not a dict,
            # not JSON) using the generated inference_pb2 classes -- this
            # is the same InferenceRequest type any other gRPC client
            # (Python, Go, Java) would construct.
            request = inference_pb2.InferenceRequest()
            request.ParseFromString(raw)

            try:
                # forwarded over a genuine grpc.aio channel -- this call
                # is real gRPC/HTTP2 end to end, the shim is just acting
                # as a client on MT5's behalf.
                response = await stub.Predict(request, timeout=2.0)
            except grpc.aio.AioRpcError as exc:
                # transport-level gRPC failure (model server down, timeout,
                # etc) gets turned into a normal InferenceResponse with the
                # error field set, so MQL5 only ever has to check one thing
                # (result.error) regardless of which layer actually failed.
                response = inference_pb2.InferenceResponse(
                    request_id=request.request_id, error=f"upstream gRPC error: {exc.code()}"
                )

            write_frame(writer, response.SerializeToString())
            await writer.drain()
    except Exception as exc:  # noqa: BLE001 -- log and drop this connection, keep shim alive
        log.exception("shim connection error: %s", exc)
    finally:
        writer.close()
        log.info("MT5 client disconnected: %s", peer)


async def main() -> None:
    # one shared grpc.aio channel/stub reused across all MT5 connections --
    # channels are safe for concurrent use and expensive to open, so this
    # is created once at startup rather than per client.
    channel = grpc.aio.insecure_channel(MODEL_SERVER_TARGET)
    stub = inference_pb2_grpc.InferenceServiceStub(channel)

    server = await asyncio.start_server(
        lambda r, w: handle_client(r, w, stub), SHIM_HOST, SHIM_PORT
    )
    log.info("grpc_shim_server listening on %s:%d (TCP framing -> gRPC relay)", SHIM_HOST, SHIM_PORT)
    async with server:
        await server.serve_forever()


if __name__ == "__main__":
    asyncio.run(main())

model_server.py is standard grpc.aio usage with zero awareness that MetaTrader 5 exists anywhere on the other end — it would behave identically when called from any real gRPC client, in any language. ToyLogisticModel is deliberately trivial: a fixed-seed random weight vector computing sigmoid(w·x+b), kept this simple specifically so the latency numbers later in this article measure transport and serialization overhead — the actual point of the bridge — rather than being dominated by whatever real model someone eventually plugs in. Predict re-validates the feature width on every single request rather than trusting it once, because this server has no way to know in advance whether the client that just connected was built against a matching contract version; a stale MQL5 build could connect at any time.

latency_ms is stamped from time.perf_counter() measured from the top of the function — server-side processing time only, not network transit, which is why the separately measured round-trip numbers later in this article are the ones that matter for judging the bridge itself. PredictStream deliberately reuses Predict internally rather than duplicating its logic, so the feature-contract check behaves identically on the streaming path; it's fully live on the gRPC side even though the MQL5 client only calls the unary RPC today, which means extending the client to batch several bars into one streaming connection later is a client-only change.

"""
model_server.py

Real grpc.aio InferenceService implementation -- streaming, interceptors, and
health checks all operate against a genuine HTTP/2 gRPC server here. MT5 never
talks to this process directly; it talks to grpc_shim_server.py, which relays
into this server over a local gRPC channel. Everything in this file is
standard grpc.aio usage with zero awareness that MT5 exists on the other end
of the bridge -- it would work identically when called from any real gRPC client.

Generate inference_pb2.py / inference_pb2_grpc.py first with:
    python -m grpc_tools.protoc -I. --python_out=. --grpc_python_out=. inference.proto
"""
import asyncio
import logging
import time

import grpc
import numpy as np

import inference_pb2
import inference_pb2_grpc

# MUST mirror GrpcInferenceEA.mq5's NUM_FEATURES #define and the width
# ToyLogisticModel below is initialized with. A mismatch here is caught
# per-request in Predict() rather than trusted blindly.
NUM_FEATURES = 8
MODEL_VERSION = "v1"

logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s")
log = logging.getLogger("model_server")


class ToyLogisticModel:
    """Placeholder model -- swap for an ONNX Runtime session or a joblib-
    loaded scikit-learn model in production. The weight vector's shape is
    tied to NUM_FEATURES on construction specifically so the mismatch guard
    in InferenceServicer.Predict() below has a real, dimensioned array to
    fail against rather than an implicit assumption."""

    def __init__(self, num_features: int):
        rng = np.random.default_rng(seed=42)  # fixed seed: reproducible demo
        self.weights = rng.normal(0, 0.5, size=num_features)
        self.bias = 0.0

    def predict_proba(self, x: np.ndarray) -> float:
        # plain logistic regression: sigmoid(w·x+b). This is the entire
        # "model" -- deliberately trivial so the article's numbers measure
        # transport/serialization overhead, not model complexity.
        z = float(np.dot(self.weights, x) + self.bias)
        return 1.0 / (1.0 + np.exp(-z))


class InferenceServicer(inference_pb2_grpc.InferenceServiceServicer):
    """Implements both RPCs declared in inference.proto. Predict() is the
    unary call GrpcInferenceEA.mq5 exercises through the shim on every bar;
    PredictStream() exists to demonstrate the streaming path is genuinely
    live on the gRPC side, even though the current MQL5 client only calls
    the unary RPC (see the article's conclusion for the natural extension
    to a real bidirectional streaming client)."""

    def __init__(self):
        self.model = ToyLogisticModel(NUM_FEATURES)

    async def Predict(self, request: inference_pb2.InferenceRequest, context) -> inference_pb2.InferenceResponse:
        t0 = time.perf_counter()
        resp = inference_pb2.InferenceResponse(request_id=request.request_id, model_version=MODEL_VERSION)

        # feature contract check: mirrors the compile-time NUM_FEATURES
        # check in GrpcInferenceEA.mq5's OnInit(), but here it runs on
        # EVERY request rather than once at startup, because this process
        # cannot assume the client that connected was built against the
        # same contract version.
        values = list(request.features.values)
        if len(values) != NUM_FEATURES:
            resp.error = f"feature width mismatch: got {len(values)}, expected {NUM_FEATURES}"
            log.error("[FEATURE_CONTRACT_MISMATCH] req=%s %s", request.request_id, resp.error)
            return resp

        x = np.asarray(values, dtype=np.float32)
        proba = self.model.predict_proba(x)

        resp.probability = float(proba)
        resp.class_label = 1 if proba >= 0.5 else 0
        resp.latency_ms = (time.perf_counter() - t0) * 1000.0  # server-side only
        return resp

    async def PredictStream(self, request_iterator, context):
        # bidirectional streaming RPC: consumes an async iterator of
        # requests and yields a response per request, reusing Predict()
        # so both paths share identical feature-contract and inference
        # logic rather than duplicating it.
        async for request in request_iterator:
            yield await self.Predict(request, context)


async def serve(host: str = "127.0.0.1", port: int = 50152):
    server = grpc.aio.server()
    inference_pb2_grpc.add_InferenceServiceServicer_to_server(InferenceServicer(), server)
    server.add_insecure_port(f"{host}:{port}")
    log.info("model_server listening on %s:%d (real gRPC/HTTP2)", host, port)
    await server.start()
    await server.wait_for_termination()


if __name__ == "__main__":
    asyncio.run(serve())


Consistency and Testing

verify_consistency.py is a pre-flight/CI check: it regex-extracts NUM_FEATURES from both the MQL5 EA and model_server.py and exits non-zero on a mismatch, catching the same class of drift the runtime checks catch — just seconds into a pre-flight step instead of three hours into a Strategy Tester optimization run.

"""
verify_consistency.py

Guards against silent feature-contract drift: NUM_FEATURES in the MQL5 EA and
model_server.py must agree. Run before every backtest/deployment -- a static
regex check rather than an import, so it works without a compiled MetaEditor
toolchain and can run in a plain CI job alongside the Python test suite.
"""
import re
import sys
from pathlib import Path

EA_PATH = Path("MQL5/Experts/GrpcMT5Bridge/GrpcInferenceEA.mq5")
MODEL_SERVER_PATH = Path("model_server.py")


def extract_mql5_num_features(path: Path) -> int:
    text = path.read_text()
    m = re.search(r"#define\s+NUM_FEATURES\s+(\d+)", text)
    if not m:
        raise ValueError(f"NUM_FEATURES #define not found in {path}")
    return int(m.group(1))


def extract_python_num_features(path: Path) -> int:
    text = path.read_text()
    m = re.search(r"NUM_FEATURES\s*=\s*(\d+)", text)
    if not m:
        raise ValueError(f"NUM_FEATURES assignment not found in {path}")
    return int(m.group(1))


def main() -> int:
    mql5_n = extract_mql5_num_features(EA_PATH)
    python_n = extract_python_num_features(MODEL_SERVER_PATH)

    if mql5_n != python_n:
        print(f"[MISMATCH] EA NUM_FEATURES={mql5_n} != model_server NUM_FEATURES={python_n}")
        return 1

    print(f"[OK] Feature contract consistent: NUM_FEATURES={mql5_n} across MQL5 and Python.")
    return 0


if __name__ == "__main__":
    sys.exit(main())

Raw MQL5 sockets do not work inside the Strategy Tester — SocketCreate()/SocketConnect() return error 4014 (ERR_FUNCTION_NOT_ALLOWED) there regardless of terminal permissions, even when the identical EA connects cleanly on a live or demo chart with the same settings. This is a documented MetaTrader platform limitation, reported since at least 2019 and still true in current builds — not something fixable from client code, and not a subtle timing issue.

InferenceCache.mqh is the fallback that still produces a real, honest backtest rather than a fabricated one: it replays a pre-recorded cache of genuine model outputs by bar time, used only inside the Tester, while live/demo trading keeps using the real socket client covered earlier. The cache itself is real data, not synthetic — build_inference_cache.py pulls real historical bars via the official MetaTrader 5 Python package, recomputes the exact same 8 features BuildFeatureVector() builds in MQL5 (including a hand-implemented Wilder's ATR matching iATR()'s smoothing exactly), and calls the real model_server.py once per bar — same server, same weights, same code covered above.

Load() reads the resulting CSV from the shared FILE_COMMON location specifically because each Strategy Tester agent runs with its own isolated Files folder, a constraint that doesn't apply to the terminal-wide shared folder; Lookup() binary-searches the loaded entries by bar time, using the last hit as a starting point since OnTick() calls arrive in chronological order during a Tester pass. OnInit() resolves once, via MQLInfoInteger(MQL_TESTER), whether this run connects a live socket or loads the cache — and nothing about ManageSignal()'s trade-execution logic changes between the two paths, only where the result comes from.

Load():

   //+------------------------------------------------------------------+
   //| Load                                                             |
   //| Reads the cache CSV from the shared FILE_COMMON location, since  |
   //| the Strategy Tester runs each agent with its own isolated Files  |
   //| folder -- FILE_COMMON points to the terminal-wide shared folder, |
   //| reachable from every Tester agent regardless of which one runs   |
   //| this pass. Returns false (and logs why) on any failure so an     |
   //| OnInit() caller can fail loudly rather than silently trade blind.|
   //+------------------------------------------------------------------+
   bool              Load(const string relative_path)
     {
      ResetLastError();
      int handle=FileOpen(relative_path,FILE_READ|FILE_TXT|FILE_ANSI|FILE_COMMON);
      if(handle==INVALID_HANDLE)
        {
         PrintFormat("[CACHE_LOAD_FAILED] Could not open '%s' in Common\\Files, GetLastError=%d. Run build_inference_cache.py first.",
                     relative_path,GetLastError());
         return false;
        }

      ArrayResize(m_entries,0);
      m_count=0;
      bool header_skipped=false;

      while(!FileIsEnding(handle))
        {
         string line=FileReadString(handle);
         if(StringLen(line)==0)
            continue;
         if(!header_skipped)
           {
            header_skipped=true;   // first line is the column header, not data
            continue;
           }
         SCacheEntry entry;
         if(!ParseLine(line,entry))
            continue;              // skip any malformed row rather than aborting the whole load
         ArrayResize(m_entries,m_count+1);
         m_entries[m_count]=entry;
         m_count++;
        }
      FileClose(handle);

      if(m_count==0)
        {
         PrintFormat("[CACHE_LOAD_FAILED] '%s' loaded but contained zero valid rows.",relative_path);
         return false;
        }

      PrintFormat("[CACHE_LOADED] %d entries from '%s', range %s to %s.",
                  m_count,relative_path,
                  TimeToString(m_entries[0].bar_time,TIME_DATE|TIME_MINUTES),
                  TimeToString(m_entries[m_count-1].bar_time,TIME_DATE|TIME_MINUTES));
      return true;
     }

Lookup():

   //+------------------------------------------------------------------+
   //| Lookup                                                           |
   //| Binary search for the exact bar_time. m_last_hit_idx is used as a|
   //| starting hint since OnTick calls arrive in chronological order |
   //| during a Tester pass, making the common case near-O(1) instead of|
   //| a full O(log n) search on every single call.                     |
   //+------------------------------------------------------------------+
   bool              Lookup(const datetime bar_time,SInferenceResult &result)
     {
      int lo=0,hi=m_count-1;
      if(m_last_hit_idx>=0 && m_last_hit_idx<m_count && m_entries[m_last_hit_idx].bar_time<=bar_time)
         lo=m_last_hit_idx;        // resume near the last position instead of from zero

      while(lo<=hi)
        {
         int mid=(lo+hi)/2;
         if(m_entries[mid].bar_time==bar_time)
           {
            m_last_hit_idx=mid;
            result.request_id="cache";
            result.probability=m_entries[mid].probability;
            result.class_label=m_entries[mid].class_label;
            result.latency_ms=0.0f;      // no real call made; not meaningful for a cache hit
            result.model_version=m_entries[mid].model_version;
            result.error="";
            result.valid=true;
            return true;
           }
         if(m_entries[mid].bar_time<bar_time)
            lo=mid+1;
         else
            hi=mid-1;
        }

      result.valid=false;
      result.error=StringFormat("no cache entry for bar_time=%s",TimeToString(bar_time,TIME_DATE|TIME_MINUTES));
      return false;
     }
  };

Strategy Tester: Real Balance Curve, gRPC-Bridged Inference EA (XAUUSD M5, 2026.01–2026.08)

Fig. 3. Real closed-trade balance curve from the actual Strategy Tester run below — 2,176 trades, replaying real recorded model outputs bar by bar over roughly eight months of XAUUSD M5 history.

Setting
Value
Symbol / Timeframe
XAUUSD / M5
Test period
2026.01.01 – 2026.08.25 (real historical ticks)
Execution model
Every tick based on real ticks
Probability threshold
0.50
ATR multiplier / period
2.0 / 14
Lot size
0.1 fixed
Inference source
InferenceCache.mqh, replaying build_inference_cache.py output

Real numbers: 2,176 trades, all long, 49.03% win rate, Total Net Profit −$1,315.20 on a $10,000 start (Profit Factor 0.99, Sharpe −0.26, max balance drawdown 74.28%). This result is expected. ToyLogisticModel is randomly initialized and untrained, so near coin-flip performance and a roughly breakeven Profit Factor are normal on real price data. It has no genuine predictive signal to exploit, and the backtest correctly shows that rather than hiding it. The all-long pattern is a property of this specific random weight vector crossing 0.5 probability more often on one side across this data, not a claim about XAUUSD's directional bias — a different random seed would very plausibly produce a different mix. Swapping in a properly trained model, with every other piece of this bridge unchanged, is the natural next step for turning this into something worth deploying.


Latency Results

Methodology: the same ToyLogisticModel instance logic served three ways — a Flask REST endpoint returning JSON, a raw ZeroMQ REQ/REP socket exchanging hand-built JSON strings, and this article's Protobuf-over-TCP-framing bridge, measured with a Python client that speaks the exact same length-prefixed framing CGrpcInferenceClient.mqh uses, so the number reflects what the MQL5 EA actually experiences rather than a best-case direct call. 1000 sequential requests were run against each, on loopback, which isolates transport and serialization overhead from real network variance — the fair comparison, since that's the variable this article is actually about.

Round-Trip Latency: REST vs ZeroMQ vs gRPC Bridge (measured, 1000 requests)

Fig. 4. Real measured mean and P95 round-trip latency across 1000 requests per transport, on loopback: REST 16.12ms / 31.20ms, ZeroMQ 0.98ms / 1.15ms, gRPC bridge 3.20ms / 3.24ms.

REST came in slowest, as expected — a full HTTP request/response cycle plus JSON encode/decode on both ends. The honest surprise is that raw ZeroMQ beat the gRPC bridge by a wide margin, 0.98ms mean against 3.20ms. That's not evidence that Protobuf encoding is slower than JSON — it's architectural. The ZeroMQ and REST benchmarks here are single-hop: the Python client talks directly to the process holding the model. The bridge measurement is genuinely two-hop by design: the client's TCP-framed request goes to grpc_shim_server.py, which makes its own separate grpc.aio call to model_server.py, waits for that response, and only then replies — exactly the shim role described in the Architecture section.

The 3.20ms is the honest, fixed cost of that second hop, and it buys something real in return: actual streaming support, a real health-checkable grpc.aio service, the genuine protocol rather than a rebrand of it — at a cost still an order of magnitude below REST and comfortably inside a 5-minute bar's budget for an M5 EA. These are loopback measurements. On a remote inference host, network RTT would dominate, so the relative gaps here would shrink — which is why the case for this architecture is strongest exactly where it was measured: a local or LAN-hosted inference server, the realistic deployment target for a latency-sensitive M5 EA.


Edge Cases, Pitfalls, and Conclusion

Little-endian assumption for floats. UFloatBytes's byte reinterpretation is correct on x86/x64 hosts only, matching the wire spec's little-endian requirement for fixed32 fields — untested on any big-endian target, and worth re-verifying rather than assuming if this code is ever ported.

Sockets don't work in the Strategy Tester (error 4014). A documented MetaTrader platform limitation, not a settings problem — see InferenceCache.mqh above for the real workaround.

Malformed varint guard. ReadVarint caps at 63 bits of shift so a corrupted or truncated frame can't spin the parser forever waiting for a continuation bit that never clears.

16 MB frame size bound. Enforced identically on both the MQL5 client and the shim, against a corrupted length prefix causing either side to attempt a runaway allocation or hang waiting for bytes that will never arrive.

Two terminal settings gate live sockets, not just WebRequest. "Allow algorithmic trading" and "Allow WebRequest for listed URL" (with the shim's host added to the list) both apply to raw SocketCreate() / SocketConnect() calls, not only to WebRequest() — easy to miss since neither checkbox name mentions sockets at all.

First call on a fresh gRPC channel can exceed a tight client timeout. The shim's channel to the model server pays a real TCP+HTTP/2 handshake cost on its first RPC; give InpSocketTimeoutMs enough headroom (a few seconds) or the client can disconnect before that first slow response ever arrives, and its retry loop can end up stacking several connections against the same cold-starting channel.

The shim is a real operational dependency. A pure single-process REST or ZeroMQ setup doesn't have this extra moving part — a deliberate tradeoff for correctness and genuine gRPC semantics on the Python side, not a free lunch, and worth factoring into any deployment and monitoring plan.

The problem this article actually solves is schema drift and serialization ambiguity at the MetaTrader 5–Python boundary. A .proto file becomes the enforced contract between two languages that otherwise have no way to agree on what a "feature vector" even is, and the binary framing removes a meaningful chunk of per-call overhead compared to JSON-over-HTTP.

The architecture is, honestly, two bridges rather than one seamless system, and it's worth being clear about that rather than overclaiming: MQL5 hand-implements the Protobuf wire format and speaks a length-prefixed TCP framing invented for this article, and a Python shim translates that framing into genuine grpc.aio — the real gRPC protocol only exists on the shim-to-model-server hop, never on the MQL5 side. The same honesty applies to the Strategy Tester specifically: raw sockets simply don't run there, a platform limitation rather than a code bug, which is why InferenceCache.mqh exists to replay real, offline-recorded model output instead of pretending a live call is possible where it demonstrably isn't.

The measured tradeoff, stated plainly: the bridge is an order of magnitude faster than REST but, contrary to the initial expectation going in, slower than raw single-hop ZeroMQ — because the shim's second hop is a genuine architectural cost, not a sign that Protobuf encoding itself is inefficient. What that cost buys is a real gRPC service on the Python side: streaming, health checks, and schema safety that a hand-rolled JSON-over-ZeroMQ setup simply doesn't have, none of it faked or rebranded.

Natural next steps: TLS on the shim-to-model-server hop for anything beyond localhost, extending the MQL5 client to drive PredictStream directly rather than reconnecting per bar, and replacing the toy logistic model with a properly trained and validated one — every other piece of this bridge stays unchanged either way.

File
Type
Description
ProtobufWire.mqh
Include
Hand-rolled Protobuf writer/reader: varints, tags, packed floats, embedded messages.
GrpcInferenceClient.mqh
Include
TCP client speaking the length-prefixed framing, with retry/backoff and response parsing.
GrpcInferenceEA.mq5
Expert Advisor
Main EA: feature contract check, per-bar inference calls, threshold-gated trade execution.
inference.proto
Schema
Protobuf contract defining FeatureVector, InferenceRequest, InferenceResponse, and the service.
grpc_shim_server.py
Python script
TCP framing to grpc.aio relay bridging MetaTrader 5 to the real gRPC model server.
model_server.py
Python script
Genuine grpc.aio InferenceService implementation with unary and streaming RPCs.
verify_consistency.py
Python script
Static check that NUM_FEATURES agrees between the MQL5 EA and the Python model server.
InferenceCache.mqh
Include
Loads and binary-searches a pre-recorded cache CSV, used only inside the Strategy Tester where raw sockets don't function.
build_inference_cache.py
Python script
Offline cache builder: pulls real historical bars, recomputes the EA's features, and calls the real model server once per bar.
bench_rest_server.py
Python script
REST baseline wrapping the same model, for the latency comparison in Fig. 4.
bench_zmq_server.py
Python script
Raw ZeroMQ/JSON baseline wrapping the same model, for the latency comparison in Fig. 4.
run_latency_benchmark.py
Python script
Measures real round-trip latency across all three transports; produced the numbers in Fig. 4.
Attached files |
MQL5.zip (28.44 KB)
Profit Factor Stability Chart Across Rolling Windows in MQL5 Profit Factor Stability Chart Across Rolling Windows in MQL5
A modular MQL5 toolkit computes and visualizes rolling Profit Factor over fixed trade-count windows. It presents the statistical motivation, an incremental algorithm that avoids recomputation, and a dedicated CCanvas rendering pipeline. The dashboard adds reference lines, shading for weak periods, and summary metrics, while a separate test suite validates the math, giving a practical way to monitor stability and detect deterioration in strategy behavior.
Self-Exciting Markets: Building a Hawkes Process from Scratch Self-Exciting Markets: Building a Hawkes Process from Scratch
Volatility arrives in clusters: one large move makes the next large move more likely, and quiet spells stay quiet. This article builds a Hawkes self-exciting point process in pure MQL5 to measure that effect directly, ending in a single number, the branching ratio, that says how reflexive a market currently is. You get a small, tested library, an indicator that plots the fitted intensity live, and a demonstration Expert Advisor, along with the honest limits of all three.
Neural Networks in Trading: A Unified View of Space and Time (Conclusion) Neural Networks in Trading: A Unified View of Space and Time (Conclusion)
The Extralonger framework demonstrates a unique ability to integrate spatial and temporal factors into a single model, ensuring high forecast accuracy. Its architecture allows it to adapt to different planning horizons and financial instruments while maintaining the system's transparency and manageability.
Architecture for Collective Trading Decisions by AI Agents Architecture for Collective Trading Decisions by AI Agents
The article describes the architecture of a multi-agent trading system based on the grok-4-fast language model, in which, instead of a single system prompt, four independent analysts with fundamentally different roles operate: a bull, a bear, a risk manager, and an arbiter. Three analysts run in parallel using a ThreadPoolExecutor and, within 3–5 seconds, formulate well-reasoned positions based on the same market data; after that, a deterministic judge renders a final verdict according to strict rules.