306 lines
No EOL
10 KiB
MQL5
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; }
|
|
}; |