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