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