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

111 lines
4.5 KiB
MQL5

//+------------------------------------------------------------------+
//| 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 */ }
}
//+------------------------------------------------------------------+