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