Article-24929-ZMQ-SUB/mock_mempool_fee_publisher.py

177 lines
5.8 KiB
Python

2026-09-23 12:31:36 +00:00
import time
import json
import random
import logging
import statistics
from collections import deque
from dataclasses import dataclass, field
from typing import Callable
import zmq
# --------------------------------------------------------------------------
# Configuration
# --------------------------------------------------------------------------
ZMQ_BIND_ADDR = "tcp://*:5556"
SCHEDULER_TICK_SECONDS = 1
DEFAULT_POLL_INTERVAL_SECONDS = 1
ROLLING_WINDOW_SIZE = 60
HEARTBEAT_TOPIC = "onchain.heartbeat"
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(message)s",
)
log = logging.getLogger("mempool_fee_publisher_mock")
# --------------------------------------------------------------------------
# Rolling z-score tracking
# --------------------------------------------------------------------------
@dataclass
class RollingStat:
window: deque = field(default_factory=lambda: deque(maxlen=ROLLING_WINDOW_SIZE))
def update_and_score(self, value: float) -> float:
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)
# --------------------------------------------------------------------------
# Mock Data Generators (Placeholders)
# --------------------------------------------------------------------------
def fetch_btc_fee_metrics_mock() -> list[dict]:
"""Generates random BTC fee rates (sat/vByte)."""
base_fee = random.uniform(10.0, 50.0)
return [
{"asset": "BTC", "metric": "fee_fastest", "value": round(base_fee * 1.5, 2)},
{"asset": "BTC", "metric": "fee_halfhour", "value": round(base_fee * 1.2, 2)},
{"asset": "BTC", "metric": "fee_hour", "value": round(base_fee, 2)},
{"asset": "BTC", "metric": "fee_economy", "value": round(base_fee * 0.8, 2)},
]
def fetch_btc_mempool_state_mock() -> list[dict]:
"""Generates random BTC mempool congestion values."""
return [
{"asset": "BTC", "metric": "mempool_pending_count", "value": float(random.randint(5000, 50000))},
{"asset": "BTC", "metric": "mempool_vsize", "value": float(random.randint(10000000, 80000000))},
]
def fetch_eth_gas_metrics_mock() -> list[dict]:
"""Generates random ETH gas prices (Gwei)."""
base_gas = random.uniform(15.0, 80.0)
return [
{"asset": "ETH", "metric": "gas_safe", "value": round(base_gas * 0.9, 2)},
{"asset": "ETH", "metric": "gas_propose", "value": round(base_gas, 2)},
{"asset": "ETH", "metric": "gas_fast", "value": round(base_gas * 1.3, 2)},
]
# --------------------------------------------------------------------------
# Fetcher Scheduling
# --------------------------------------------------------------------------
@dataclass
class ScheduledFetcher:
name: str
fetch_fn: Callable[[], list[dict]]
poll_interval_sec: float = DEFAULT_POLL_INTERVAL_SECONDS
next_due: float = 0.0
def is_due(self, now: float) -> bool:
return now >= self.next_due
def on_success(self, now: float) -> None:
self.next_due = now + self.poll_interval_sec
FETCHERS: list[ScheduledFetcher] = [
ScheduledFetcher("BTC-fees-mock", fetch_btc_fee_metrics_mock, poll_interval_sec=2),
ScheduledFetcher("BTC-mempool-mock", fetch_btc_mempool_state_mock, poll_interval_sec=2),
ScheduledFetcher("ETH-gas-mock", fetch_eth_gas_metrics_mock, poll_interval_sec=2)
]
# --------------------------------------------------------------------------
# ZMQ Publishing & Main Loop
# --------------------------------------------------------------------------
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:
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)
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
raw_records = fetcher.fetch_fn()
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 broadcasting mock feeds bound at %s", ZMQ_BIND_ADDR)
time.sleep(1.0)
while True:
run_due_fetchers(pub)
publish_heartbeat(pub)
time.sleep(SCHEDULER_TICK_SECONDS)
if __name__ == "__main__":
try:
main()
except KeyboardInterrupt:
log.info("shutting down")