303 lines
11 KiB
Python
303 lines
11 KiB
Python
"""
| |||
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")
|