//+------------------------------------------------------------------+ //| OnChainFeedSubscriber.mq5 | //| Copyright 2025, MetaQuotes Ltd. | //| https://www.mql5.com | //| Minimal demonstration of COnchainFeed (OnchainFeed.mqh) wired | //| into an EA's lifecycle. Subscribes to the mempool-fee topics | //| published by mempool_fee_publisher.py and logs the cached state | //| on a timer. | //+------------------------------------------------------------------+ #property copyright "Copyright 2025, MetaQuotes Ltd." #property link "https://www.mql5.com" #property version "1.00" #include "OnchainFeed.mqh" //--- input parameters input string InpPublisherHost = "LOCALHOST"; input int InpPublisherPort = 5556; input int InpPollIntervalMs = 500; input int InpHeartbeatMaxAgeSec = 50; input int InpMetricMaxAgeSec = 300; //--- Select what feed(s) to subscribe to input bool InpSubscribeToBTC_fastestfee = true; input bool InpSubscribeToBTC_mempoolpendingcount = false; input bool InpSubscribeToETH_gasfast = false; //--- COnchainFeed g_feed; //+------------------------------------------------------------------+ int OnInit() { // -- Make sure at least one subscription has been enabled. if(!InpSubscribeToBTC_fastestfee && !InpSubscribeToBTC_mempoolpendingcount && !InpSubscribeToETH_gasfast) { Alert("You have not subscribed to any feeds."); return INIT_FAILED; } g_feed.SetStalenessThresholds(InpHeartbeatMaxAgeSec, InpMetricMaxAgeSec); // Subscribe to the specific fee/congestion topics this EA cares // about, across multiple assets. ZMQ SUB filtering is prefix-based, // so each of these only matches its own topic -- we are not // subscribing to "" (everything). if(InpSubscribeToBTC_fastestfee) g_feed.AddSubscription("onchain.BTC.fee_fastest"); if(InpSubscribeToBTC_mempoolpendingcount) g_feed.AddSubscription("onchain.BTC.mempool_pending_count"); if(InpSubscribeToETH_gasfast) g_feed.AddSubscription("onchain.ETH.gas_fast"); //--- g_feed.AddSubscription("onchain.heartbeat"); //--- if(!g_feed.Connect(InpPublisherHost, InpPublisherPort)) Print("OnchainFeedDemoEA: initial connect failed, will retry via timer"); //--- EventSetMillisecondTimer(InpPollIntervalMs); return(INIT_SUCCEEDED); } //+------------------------------------------------------------------+ void OnDeinit(const int reason) { Comment(""); EventKillTimer(); } //+------------------------------------------------------------------+ // All socket I/O happens here, not in OnTick() -- keeps tick handling // free of network calls regardless of how often ticks arrive. void OnTimer() { g_feed.Poll(); } //+------------------------------------------------------------------+ void OnTick() { CheckOnchainSignal(); } //+------------------------------------------------------------------+ // Example gating logic: both staleness checks must pass before the // cached value is trusted for anything. void CheckOnchainSignal() { OnchainMetric btc_fee, eth_fee,mem_count; bool have_btc = g_feed.GetMetric("BTC.fee_fastest", btc_fee); bool have_eth = g_feed.GetMetric("ETH.gas_fast", eth_fee); bool have_mempoolcount = g_feed.GetMetric("BTC.mempool_pending_count",mem_count); if(!have_btc && !have_eth && !have_mempoolcount) return; // no messages received for either topic yet since startup if(!g_feed.IsFeedAlive()) { Print("Onchain feed: STALE (no heartbeat within threshold)"); return; } string status = ""; if(have_btc && g_feed.IsMetricFresh(btc_fee)) status += StringFormat("BTC fee_fastest=%.1f (z=%.2f) ", btc_fee.value, btc_fee.zscore); if(have_mempoolcount && g_feed.IsMetricFresh(mem_count)) status += StringFormat("BTC mempool pending count=%.1f (z=%.2f) ", mem_count.value,mem_count.zscore); if(have_eth && g_feed.IsMetricFresh(eth_fee)) status += StringFormat("ETH gas_fast=%.1f (z=%.2f) ", eth_fee.value, eth_fee.zscore); Comment("Onchain feed: " + (StringLen(status) > 0 ? status : "no fresh metrics")); // Do something with feed data // if(have_btc && btc_fee.zscore > 2.0) { /* elevated BTC congestion */ } } //+------------------------------------------------------------------+