Article-24929-ZMQ-SUB/mock_mempool_fee_publisher.py
2026-09-23 12:31:36 +00:00

177 lines
No EOL
5.8 KiB
Python

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