Article-24929-ZMQ-SUB/mempool_fee_publisher.py

303 lines
11 KiB
Python

2026-09-23 12:31:36 +00:00
"""
mempool_fee_publisher.py
Fetches on-chain fee/congestion metrics from multiple free, no-label-required
sources and publishes them over a ZMQ PUB socket for MQL5 SUB clients (via
Zmtp.mqh) to consume.
Topics published (multi-part ZMQ messages: [topic, json_payload]):
onchain.BTC.fee_fastest
onchain.BTC.fee_halfhour
onchain.BTC.fee_hour
onchain.BTC.fee_economy
onchain.BTC.mempool_pending_count
onchain.BTC.mempool_vsize
onchain.ETH.gas_safe
onchain.ETH.gas_propose
onchain.ETH.gas_fast
onchain.heartbeat <- liveness signal, published on every
scheduler tick regardless of whether
any individual source succeeded
The heartbeat lets a SUB client distinguish "the publisher process died"
from "one particular source is having trouble" -- it fires unconditionally
each tick, separate from the per-metric timestamps above.
Payload schema (flat, MQL5-friendly):
{
"asset": "BTC",
"metric": "fee_fastest",
"value": 42.0,
"zscore": 1.83,
"timestamp": 1725500000,
"window": "288" # number of samples in the rolling window
}
"""
import os
import time
import json
import logging
import statistics
from collections import deque
from dataclasses import dataclass, field
from typing import Callable
import requests
import zmq
# --------------------------------------------------------------------------
# Configuration
# --------------------------------------------------------------------------
ZMQ_BIND_ADDR = "tcp://*:5556"
SCHEDULER_TICK_SECONDS = 1 # how often the main loop wakes up to check
# each fetcher's due time -- NOT the fetch
# interval itself
DEFAULT_POLL_INTERVAL_SECONDS = 1
ROLLING_WINDOW_SIZE = 60 # samples of history for z-score baseline
HTTP_TIMEOUT_SECONDS = 5
MAX_BACKOFF_SECONDS = 300
MEMPOOL_FEES_URL = "https://mempool.fractalbitcoin.io/api/v1/fees/recommended"
MEMPOOL_STATE_URL = "https://mempool.fractalbitcoin.io/api/mempool"
ETHERSCAN_API_KEY = os.environ.get("ETHERSCAN_API_KEY", "YourApiKeyToken")
ETHERSCAN_GAS_ORACLE_URL = "https://api.etherscan.io/v2/api"
HEARTBEAT_TOPIC = "onchain.heartbeat"
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(message)s",
)
log = logging.getLogger("mempool_fee_publisher")
if ETHERSCAN_API_KEY == "YourApiKeyToken":
log.warning(
"ETHERSCAN_API_KEY not set -- using the shared demo token, which is "
"heavily rate-limited. Set the ETHERSCAN_API_KEY env var with a free "
"key from https://etherscan.io/apis for reliable ETH gas data."
)
# --------------------------------------------------------------------------
# Rolling z-score tracking, one per (asset, metric) pair
# --------------------------------------------------------------------------
@dataclass
class RollingStat:
window: deque = field(default_factory=lambda: deque(maxlen=ROLLING_WINDOW_SIZE))
def update_and_score(self, value: float) -> float:
"""Add value to the window, return its z-score against prior history.
Returns 0.0 until enough samples exist to compute a meaningful stdev,
so the EA never sees a wild early-run z-score off 1-2 data points.
"""
if len(self.window) < 10:
self.window.append(value)
return 0.0
mean = statistics.mean(self.window)
stdev = statistics.pstdev(self.window)
self.window.append(value)
if stdev == 0:
return 0.0
return round((value - mean) / stdev, 3)
_stats: dict[str, RollingStat] = {}
def zscore_for(asset: str, metric: str, value: float) -> float:
key = f"{asset}.{metric}"
if key not in _stats:
_stats[key] = RollingStat()
return _stats[key].update_and_score(value)
# --------------------------------------------------------------------------
# Data fetching -- each function returns raw records (asset/metric/value
# only; zscore + timestamp + window are filled by the
# scheduler after a successful fetch, not by the fetcher itself)
# --------------------------------------------------------------------------
def fetch_btc_fee_metrics() -> list[dict]:
"""GET /api/v1/fees/recommended -> fastestFee, halfHourFee, hourFee, economyFee (sat/vByte)."""
resp = requests.get(MEMPOOL_FEES_URL, timeout=HTTP_TIMEOUT_SECONDS)
resp.raise_for_status()
data = resp.json()
field_map = {
"fastestFee": "fee_fastest",
"halfHourFee": "fee_halfhour",
"hourFee": "fee_hour",
"economyFee": "fee_economy",
}
records = []
for api_field, metric_name in field_map.items():
if api_field in data:
records.append({"asset": "BTC", "metric": metric_name, "value": float(data[api_field])})
return records
def fetch_btc_mempool_state() -> list[dict]:
"""GET /api/mempool -> pending tx count and total vsize (congestion proxy)."""
resp = requests.get(MEMPOOL_STATE_URL, timeout=HTTP_TIMEOUT_SECONDS)
resp.raise_for_status()
data = resp.json()
records = []
if "count" in data:
records.append({"asset": "BTC", "metric": "mempool_pending_count", "value": float(data["count"])})
if "vsize" in data:
records.append({"asset": "BTC", "metric": "mempool_vsize", "value": float(data["vsize"])})
return records
def fetch_eth_gas_metrics() -> list[dict]:
"""GET Etherscan gas oracle -> Safe/Propose/Fast gas price, in Gwei."""
params = {"chainid": "1", "module": "gastracker", "action": "gasoracle", "apikey": ETHERSCAN_API_KEY}
resp = requests.get(ETHERSCAN_GAS_ORACLE_URL, params=params, timeout=HTTP_TIMEOUT_SECONDS)
resp.raise_for_status()
data = resp.json()
if data.get("status") != "1":
# Etherscan returns HTTP 200 even for rate-limit/error responses --
# the real failure signal is in the payload, not the status code.
raise RuntimeError(f"Etherscan gas oracle error: {data.get('message')} / {data.get('result')}")
result = data["result"]
field_map = {
"SafeGasPrice": "gas_safe",
"ProposeGasPrice": "gas_propose",
"FastGasPrice": "gas_fast",
}
records = []
for api_field, metric_name in field_map.items():
if api_field in result:
records.append({"asset": "ETH", "metric": metric_name, "value": float(result[api_field])})
return records
# --------------------------------------------------------------------------
# Fetcher scheduling -- each source gets its own cadence and its own
# independent backoff, so one troubled source never throttles the others
# --------------------------------------------------------------------------
@dataclass
class ScheduledFetcher:
name: str # for logging only
fetch_fn: Callable[[], list[dict]]
poll_interval_sec: float = DEFAULT_POLL_INTERVAL_SECONDS
next_due: float = 0.0 # monotonic time; 0 = due immediately
backoff_sec: float = DEFAULT_POLL_INTERVAL_SECONDS
def is_due(self, now: float) -> bool:
return now >= self.next_due
def on_success(self, now: float) -> None:
self.backoff_sec = self.poll_interval_sec # reset backoff after a clean fetch
self.next_due = now + self.poll_interval_sec
def on_failure(self, now: float) -> None:
log.warning("%s: fetch failed -- next attempt in %ss", self.name, self.backoff_sec)
self.next_due = now + self.backoff_sec
self.backoff_sec = min(self.backoff_sec * 2, MAX_BACKOFF_SECONDS)
FETCHERS: list[ScheduledFetcher] = [
ScheduledFetcher("BTC-fees", fetch_btc_fee_metrics, poll_interval_sec=2),
ScheduledFetcher("BTC-mempool", fetch_btc_mempool_state, poll_interval_sec=2),
ScheduledFetcher("ETH-gas", fetch_eth_gas_metrics, poll_interval_sec=2)
]
# --------------------------------------------------------------------------
# ZMQ publishing
# --------------------------------------------------------------------------
def publish(pub_socket: "zmq.Socket", record: dict) -> None:
topic = f"onchain.{record['asset']}.{record['metric']}"
payload = json.dumps(record)
pub_socket.send_multipart([topic.encode("utf-8"), payload.encode("utf-8")])
log.info(
"published %-32s value=%-12s zscore=%s",
topic, record["value"], record["zscore"],
)
def publish_heartbeat(pub_socket: "zmq.Socket") -> None:
"""Fire the liveness signal. Called every scheduler tick regardless of
whether any individual source's fetch succeeded -- its only job is "the
process is alive," not "the data is fresh." A SUB client should watch
this independently of any individual metric's own timestamp.
"""
payload = json.dumps({"timestamp": int(time.time() + time.localtime().tm_gmtoff)})
pub_socket.send_multipart([HEARTBEAT_TOPIC.encode("utf-8"), payload.encode("utf-8")])
log.debug("published %s", HEARTBEAT_TOPIC)
# --------------------------------------------------------------------------
# Main loop
# --------------------------------------------------------------------------
def run_due_fetchers(pub_socket: "zmq.Socket") -> None:
now = time.monotonic()
ts = int(time.time()) + time.localtime().tm_gmtoff
for fetcher in FETCHERS:
if not fetcher.is_due(now):
continue
try:
raw_records = fetcher.fetch_fn()
except requests.exceptions.RequestException as exc:
log.warning("%s: request error (%s)", fetcher.name, exc)
fetcher.on_failure(now)
continue
except Exception:
log.exception("%s: unexpected error", fetcher.name)
fetcher.on_failure(now)
continue
for raw in raw_records:
record = {
"asset": raw["asset"],
"metric": raw["metric"],
"value": raw["value"],
"zscore": zscore_for(raw["asset"], raw["metric"], raw["value"]),
"timestamp": ts,
"window": str(ROLLING_WINDOW_SIZE),
}
publish(pub_socket, record)
fetcher.on_success(now)
def main() -> None:
ctx = zmq.Context()
pub = ctx.socket(zmq.PUB)
pub.bind(ZMQ_BIND_ADDR)
log.info("ZMQ PUB socket bound at %s", ZMQ_BIND_ADDR)
log.info("scheduled fetchers: %s", ", ".join(f.name for f in FETCHERS))
time.sleep(1.0)
while True:
run_due_fetchers(pub)
publish_heartbeat(pub) # unconditional, every tick
time.sleep(SCHEDULER_TICK_SECONDS)
if __name__ == "__main__":
try:
main()
except KeyboardInterrupt:
log.info("shutting down")