177 lines
5.8 KiB
Python
177 lines
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")
|