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