Article-24929-ZMQ-SUB/OnChainFeed.mqh
2026-09-23 12:31:36 +00:00

306 lines
No EOL
10 KiB
MQL5

//+------------------------------------------------------------------+
//| OnChainFeed.mqh |
//| Copyright 2025, MetaQuotes Ltd. |
//| https://www.mql5.com |
//+------------------------------------------------------------------+
#property copyright "Copyright 2025, MetaQuotes Ltd."
#property link "https://www.mql5.com"
#include "Zmtp.mqh"
//--- one cached metric ------------------------------------------------
struct OnchainMetric
{
double value;
double zscore;
datetime timestamp; // from the publisher's payload, not local receipt time
bool valid; // false until at least one message for this key has arrived
};
//--- fixed-size metric cache, keyed by short name (e.g. "fee_fastest") -
// A handful of topics at most, so linear search over a small array is
// simpler and plenty fast compared to reaching for a hash-map pattern.
#define ONCHAIN_MAX_METRICS 32
//+------------------------------------------------------------------+
//| SUB-side cache for the onchain mempool-fee feed published by |
//| mempool_fee_publisher.py. Wraps CZmqSubSocket (Zmtp.mqh) and |
//| gives calling code a simple, non-blocking accessor. Supports |
//| multiple assets on one instance -- metrics are keyed on |
//| "<asset>.<metric>": |
//+------------------------------------------------------------------+
class COnchainFeed
{
private:
CZmqSubSocket m_sub;
string m_host;
int m_port;
int m_connect_timeout_ms; // short on purpose -- see header note
bool m_connected;
datetime m_last_reconnect_attempt;
int m_reconnect_backoff_sec;
int m_max_reconnect_backoff_sec;
string m_keys[ONCHAIN_MAX_METRICS];
OnchainMetric m_values[ONCHAIN_MAX_METRICS];
int m_count;
datetime m_last_heartbeat; // publisher-side timestamp of last heartbeat
int m_heartbeat_max_age_sec;
int m_metric_max_age_sec;
string m_subscriptions[]; // topic prefixes to (re-)subscribe on connect
//--- find or create a slot for `key` ------------------------------
int IndexOf(const string &key)
{
for(int i=0; i<m_count; i++)
if(m_keys[i] == key)
return i;
return -1;
}
int SlotFor(const string &key)
{
int idx = IndexOf(key);
if(idx >= 0)
return idx;
if(m_count >= ONCHAIN_MAX_METRICS)
{
Print("OnchainFeed: metric cache full (", ONCHAIN_MAX_METRICS, "), dropping key ", key);
return -1;
}
m_keys[m_count] = key;
m_values[m_count].valid = false;
return m_count++;
}
//--- minimal flat-JSON field extraction ---------------------------
// The publisher's payloads are deliberately flat (no nested objects/
// arrays), so a real JSON parser is overkill. This pulls the raw
// string between the colon after "key" and the next comma or
// closing brace, then lets the caller convert it.
bool ExtractRawField(const string &json, const string &key, string &out_raw)
{
string needle = "\"" + key + "\":";
int pos = StringFind(json, needle);
if(pos < 0)
return false;
int start = pos + StringLen(needle);
// skip a leading quote for string-valued fields
bool quoted = false;
if(start < StringLen(json) && StringGetCharacter(json, start) == '"')
{
quoted = true;
start++;
}
int end = start;
int len = StringLen(json);
while(end < len)
{
ushort c = StringGetCharacter(json, end);
if(quoted && c == '"')
break;
if(!quoted && (c == ',' || c == '}'))
break;
end++;
}
out_raw = StringSubstr(json, start, end - start);
return true;
}
bool ExtractDouble(const string &json, const string &key, double &out)
{
string raw;
if(!ExtractRawField(json, key, raw))
return false;
out = StringToDouble(raw);
return true;
}
bool ExtractLong(const string &json, const string &key, long &out)
{
string raw;
if(!ExtractRawField(json, key, raw))
return false;
out = StringToInteger(raw);
return true;
}
bool ExtractString(const string &json, const string &key, string &out)
{
return ExtractRawField(json, key, out);
}
//--- apply one parsed [topic, payload] message to the cache -------
void HandleMessage(const string &topic, const string &payload)
{
//Print(__FUNCTION__, "topic", topic, " payload ", payload);
if(topic == "onchain.heartbeat")
{
long ts;
if(ExtractLong(payload, "timestamp", ts))
m_last_heartbeat = (datetime)ts;
return;
}
// topic shape: onchain.<asset>.<metric> -- key the cache on
// "<asset>.<metric>" (e.g. "BTC.fee_fastest") so multiple assets
// can share one COnchainFeed instance without colliding. Callers
// pass the same "<asset>.<metric>" string to GetMetric().
string cache_key = topic;
int first_dot = StringFind(topic, ".");
if(first_dot >= 0)
cache_key = StringSubstr(topic, first_dot + 1);
int idx = SlotFor(cache_key);
if(idx < 0)
return;
double value = 0.0, zscore = 0.0;
long ts = 0;
ExtractDouble(payload, "value", value);
ExtractDouble(payload, "zscore", zscore);
ExtractLong(payload, "timestamp", ts);
m_values[idx].value = value;
m_values[idx].zscore = zscore;
m_values[idx].timestamp = (datetime)ts;
m_values[idx].valid = true;
}
public:
COnchainFeed()
{
m_connected = false;
m_count = 0;
m_last_heartbeat = 0;
m_last_reconnect_attempt = 0;
m_reconnect_backoff_sec = 1;
m_max_reconnect_backoff_sec = 60;
m_heartbeat_max_age_sec = 90; // ~3x the publisher's 30s cycle
m_metric_max_age_sec = 300; // generous vs a 30s publish cadence
m_connect_timeout_ms = 200; // see header note on the timeout tradeoff
}
// Call once before Connect(). Subscriptions are (re-)sent on every
// successful connect, since a fresh ZMTP session starts unsubscribed.
void AddSubscription(const string &topic_prefix)
{
int n = ArraySize(m_subscriptions);
ArrayResize(m_subscriptions, n + 1);
m_subscriptions[n] = topic_prefix;
}
void SetStalenessThresholds(int heartbeat_max_age_sec, int metric_max_age_sec)
{
m_heartbeat_max_age_sec = heartbeat_max_age_sec;
m_metric_max_age_sec = metric_max_age_sec;
}
bool Connect(const string &host, int port)
{
m_host = host;
m_port = port;
if(!m_sub.Connect(m_host, m_port, ZMTP_SEC_NULL, m_connect_timeout_ms))
{
Print("OnchainFeed: connect/handshake failed to ", m_host, ":", m_port);
m_connected = false;
return false;
}
for(int i = 0; i < ArraySize(m_subscriptions); i++)
m_sub.Subscribe(m_subscriptions[i]);
m_connected = true;
m_reconnect_backoff_sec = 1; // reset backoff after a clean connect
Print("OnchainFeed: connected to ", m_host, ":", m_port,
" (", ArraySize(m_subscriptions), " subscriptions)");
return true;
}
// Call from OnTimer() (recommended interval: 250ms-1s). Drains any
// backlog of pending messages, bounded by max_messages so a burst
// from the publisher can't stall this call indefinitely.
void Poll(void)
{
if(!m_connected)
{
TryReconnect();
return;
}
string topic, payload;
int drained = 0;
int max_messages = int(m_subscriptions.Size());
while(drained < max_messages && m_sub.Recv(topic, payload))
{
HandleMessage(topic, payload);
drained++;
}
// Recv() returning false means either "nothing pending within the
// connect timeout" or "the socket dropped" -- IsConnected() checks
// the socket directly, so it's the reliable way to tell them apart.
if(!m_sub.IsConnected())
{
Print("OnchainFeed: connection lost, will attempt reconnect");
m_connected = false;
}
}
void TryReconnect(void)
{
datetime now = TimeLocal();
if(now - m_last_reconnect_attempt < m_reconnect_backoff_sec)
return;
m_last_reconnect_attempt = now;
m_sub.Disconnect();
if(Connect(m_host, m_port))
return;
m_reconnect_backoff_sec = MathMin(m_reconnect_backoff_sec * 2, m_max_reconnect_backoff_sec);
Print("OnchainFeed: reconnect failed, next attempt in ", m_reconnect_backoff_sec, "s");
}
// --- accessors used by trading logic ------------------------------
// key format: "<asset>.<metric>", e.g. "BTC.fee_fastest" -- matches
// the topic with the "onchain." prefix stripped.
bool GetMetric(const string &metric_name, OnchainMetric &out)
{
int idx = IndexOf(metric_name);
if(idx < 0 || !m_values[idx].valid)
return false;
out = m_values[idx];
return true;
}
// True only if the publisher process itself is alive (heartbeat
// topic seen recently) -- independent of whether any specific
// metric has updated.
bool IsFeedAlive()
{
//Print(__FUNCTION__, "mlastheartbeat ", (int)m_last_heartbeat);
//Print(__FUNCTION__, "Timelocal ", (int)TimeLocal());
if(m_last_heartbeat == 0)
return false;
return (TimeLocal() - m_last_heartbeat) <= m_heartbeat_max_age_sec;
}
// True only if this specific metric's own timestamp is recent.
// Use alongside IsFeedAlive(), not instead of it -- see the header
// note on why both checks exist.
bool IsMetricFresh(const OnchainMetric &m)
{
if(!m.valid)
return false;
return (TimeLocal() - m.timestamp) <= m_metric_max_age_sec;
}
bool IsConnected() { return m_connected; }
};