""" 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")