111 lines
4.5 KiB
MQL5
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 */ }
|
|
}
|
|
//+------------------------------------------------------------------+
|