//+------------------------------------------------------------------+ //| 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 | //| ".": | //+------------------------------------------------------------------+ 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= 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.. -- key the cache on // "." (e.g. "BTC.fee_fastest") so multiple assets // can share one COnchainFeed instance without colliding. Callers // pass the same "." 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: ".", 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; } };