""" Maker Edge 事件记录器(Dry-run / Live)— Execution Reality Layer 事件: - quote_created / quote_canceled / quote_filled (报价生命周期) - book_tick (可选心跳,用于成交前5s盘口) - fill (成交瞬间 + 盘口状态) - fill_path (10s/30s/1m/5m + Price/Time MAE/MFE) 输出:user_data/logs/maker_edge/YYYYMMDD.jsonl """ from __future__ import annotations import json import logging import time import uuid from collections import deque from dataclasses import dataclass, field from datetime import datetime, timezone from pathlib import Path from typing import Any, Optional logger = logging.getLogger(__name__) def _utc_now() -> datetime: return datetime.now(timezone.utc) def _iso(ts: datetime | float | None = None) -> str: if ts is None: t = _utc_now() elif isinstance(ts, (int, float)): t = datetime.fromtimestamp(ts, tz=timezone.utc) else: t = ts if ts.tzinfo else ts.replace(tzinfo=timezone.utc) return t.isoformat() @dataclass class MicroSnapshot: best_bid: float = 0.0 best_ask: float = 0.0 mid: float = 0.0 spread: float = 0.0 bid_depth_1: float = 0.0 ask_depth_1: float = 0.0 bid_depth_5: float = 0.0 ask_depth_5: float = 0.0 bid_depth: float = 0.0 # top-N ask_depth: float = 0.0 obi: float = 0.0 delta: float = 0.0 trade_imbalance: float = 0.0 # (buy-sell)/(buy+sell) on recent trades delta_efficiency: float = 0.0 liquidation_distance: float = 0.0 def to_book_fields(self) -> dict[str, float]: return { "bid_price": self.best_bid, "ask_price": self.best_ask, "mid": self.mid, "spread": self.spread, "bid_depth_1": self.bid_depth_1, "ask_depth_1": self.ask_depth_1, "bid_depth_5": self.bid_depth_5, "ask_depth_5": self.ask_depth_5, "bid_depth": self.bid_depth, "ask_depth": self.ask_depth, "obi": self.obi, "delta": self.delta, "trade_imbalance": self.trade_imbalance, "delta_efficiency": self.delta_efficiency, "liquidation_distance": self.liquidation_distance, # 兼容旧字段 "buy1_depth": self.bid_depth_1, "sell1_depth": self.ask_depth_1, } @dataclass class ActiveQuote: quote_id: str pair: str side: str # bid / ask quote_price: float created_ts: float reason: str = "" trade_id: Optional[int] = None status: str = "open" # open / filled / canceled @dataclass class PendingFillPath: fill_id: str pair: str side: str fill_price: float fill_ts: float quote_id: Optional[str] = None exit_reason: Optional[str] = None # horizon prices after_10s_price: Optional[float] = None after_30s_price: Optional[float] = None after_1m_price: Optional[float] = None after_5m_price: Optional[float] = None # running extrema min_price: float = 0.0 max_price: float = 0.0 # time-MAE: worst adverse excursion seen by each horizon (signed, adverse negative for long) mae_10s: Optional[float] = None mae_30s: Optional[float] = None mae_1m: Optional[float] = None mae_5m: Optional[float] = None mfe_10s: Optional[float] = None mfe_30s: Optional[float] = None mfe_1m: Optional[float] = None mfe_5m: Optional[float] = None done: bool = False def __post_init__(self): self.min_price = self.fill_price self.max_price = self.fill_price def signed_excursions(self) -> tuple[float, float]: """Return (mae, mfe) at current min/max. mae<=0 adverse, mfe>=0 favorable.""" if self.side == "long": mae = (self.min_price - self.fill_price) / self.fill_price mfe = (self.max_price - self.fill_price) / self.fill_price else: mae = (self.fill_price - self.max_price) / self.fill_price mfe = (self.fill_price - self.min_price) / self.fill_price return mae, mfe class MakerEdgeLogger: def __init__( self, log_dir: str | Path | None = None, levels: int = 10, book_history_sec: float = 30.0, ): root = Path(__file__).resolve().parents[2] self.log_dir = Path(log_dir) if log_dir else root / "logs" / "maker_edge" self.log_dir.mkdir(parents=True, exist_ok=True) self.levels = levels self.book_history_sec = book_history_sec self._pending: dict[str, PendingFillPath] = {} self._quotes: dict[str, ActiveQuote] = {} # quote_id -> ActiveQuote self._quotes_by_trade: dict[int, str] = {} # trade_id -> quote_id self._book_hist: deque[tuple[float, MicroSnapshot]] = deque(maxlen=2000) def _file(self) -> Path: return self.log_dir / f"{_utc_now().strftime('%Y%m%d')}.jsonl" def write(self, event: dict[str, Any]) -> None: event.setdefault("ts", _iso()) event.setdefault("ts_epoch", time.time()) with self._file().open("a", encoding="utf-8") as f: f.write(json.dumps(event, ensure_ascii=False, default=str) + "\n") # ------------------------------------------------------------------ # # Snapshot # ------------------------------------------------------------------ # @staticmethod def snapshot_from_orderbook( ob: dict, levels: int = 10, recent_trades: list | None = None, last_mid: float | None = None, liq_proxy_low: float | None = None, liq_proxy_high: float | None = None, ) -> MicroSnapshot: bids = (ob.get("bids") or [])[:levels] asks = (ob.get("asks") or [])[:levels] if not bids or not asks: return MicroSnapshot() best_bid = float(bids[0][0]) best_ask = float(asks[0][0]) mid = (best_bid + best_ask) / 2.0 spread = best_ask - best_bid def depth(levels_side, n): return sum(float(x[1]) for x in levels_side[:n]) bid_depth_1 = depth(bids, 1) ask_depth_1 = depth(asks, 1) bid_depth_5 = depth(bids, 5) ask_depth_5 = depth(asks, 5) bid_depth = depth(bids, levels) ask_depth = depth(asks, levels) tot = bid_depth + ask_depth obi = ((bid_depth - ask_depth) / tot) if tot > 0 else 0.0 buy_v = sell_v = 0.0 if recent_trades: for t in recent_trades: amt = float(t.get("amount") or t.get("qty") or 0.0) side = (t.get("side") or "").lower() if side in ("buy", "b"): buy_v += amt elif side in ("sell", "s"): sell_v += amt delta = buy_v - sell_v timb_den = buy_v + sell_v trade_imbalance = ((buy_v - sell_v) / timb_den) if timb_den > 0 else 0.0 de = 0.0 if last_mid and mid and abs(delta) > 1e-12: de = ((mid - last_mid) / last_mid) / delta liq_dist = 0.0 if liq_proxy_low and liq_proxy_high and mid: rng = liq_proxy_high - liq_proxy_low if rng > 0: liq_dist = ((mid - liq_proxy_low) / rng) * 2 - 1 return MicroSnapshot( best_bid=best_bid, best_ask=best_ask, mid=mid, spread=spread, bid_depth_1=bid_depth_1, ask_depth_1=ask_depth_1, bid_depth_5=bid_depth_5, ask_depth_5=ask_depth_5, bid_depth=bid_depth, ask_depth=ask_depth, obi=obi, delta=delta, trade_imbalance=trade_imbalance, delta_efficiency=de, liquidation_distance=liq_dist, ) def record_book(self, snap: MicroSnapshot, now: float | None = None) -> None: now = now or time.time() self._book_hist.append((now, snap)) # trim old cutoff = now - self.book_history_sec while self._book_hist and self._book_hist[0][0] < cutoff: self._book_hist.popleft() def book_at(self, target_ts: float) -> Optional[MicroSnapshot]: """取最接近 target_ts 的历史盘口(用于成交前5s)。""" if not self._book_hist: return None best = min(self._book_hist, key=lambda x: abs(x[0] - target_ts)) return best[1] def book_deterioration(self, side: str, now: float | None = None, lookback: float = 5.0) -> dict: """ 成交前 lookback 秒盘口是否恶化。 long: bid_depth 下降 / ask_depth 上升 / mid 下跌 → 恶化 """ now = now or time.time() cur = self.book_at(now) past = self.book_at(now - lookback) if not cur or not past or past.mid <= 0: return {"book_ok": False} mid_chg = (cur.mid - past.mid) / past.mid bid5_chg = (cur.bid_depth_5 - past.bid_depth_5) / past.bid_depth_5 if past.bid_depth_5 else 0.0 ask5_chg = (cur.ask_depth_5 - past.ask_depth_5) / past.ask_depth_5 if past.ask_depth_5 else 0.0 obi_chg = cur.obi - past.obi if side == "long": deteriorated = (mid_chg < -0.00005) or (bid5_chg < -0.15) or (obi_chg < -0.1) else: deteriorated = (mid_chg > 0.00005) or (ask5_chg < -0.15) or (obi_chg > 0.1) return { "book_ok": True, "pre_5s_mid_chg": mid_chg, "pre_5s_bid_depth_5_chg": bid5_chg, "pre_5s_ask_depth_5_chg": ask5_chg, "pre_5s_obi_chg": obi_chg, "pre_5s_deteriorated": bool(deteriorated), "pre_5s_bid_depth_1": past.bid_depth_1, "pre_5s_ask_depth_1": past.ask_depth_1, "pre_5s_bid_depth_5": past.bid_depth_5, "pre_5s_ask_depth_5": past.ask_depth_5, "pre_5s_obi": past.obi, "pre_5s_spread": past.spread, "pre_5s_trade_imbalance": past.trade_imbalance, } # ------------------------------------------------------------------ # # Quote lifecycle # ------------------------------------------------------------------ # def create_quote( self, pair: str, side: str, quote_price: float, inventory: float, snap: MicroSnapshot, reason: str = "", trade_id: Optional[int] = None, state: dict | None = None, ) -> str: qid = uuid.uuid4().hex[:16] now = time.time() q = ActiveQuote( quote_id=qid, pair=pair, side=side, quote_price=quote_price, created_ts=now, reason=reason, trade_id=trade_id, status="open", ) self._quotes[qid] = q if trade_id is not None: self._quotes_by_trade[trade_id] = qid ev = { "event": "quote_created", "quote_id": qid, "pair": pair, "side": side, "quote_price": quote_price, "quote_created_time": _iso(now), "quote_created_epoch": now, "inventory": inventory, "reason": reason, "trade_id": trade_id, "status": "open", "filled": False, } ev.update(snap.to_book_fields()) if state: ev.update(state) self.write(ev) return qid def cancel_quote( self, quote_id: str | None = None, trade_id: Optional[int] = None, reason: str = "timeout", snap: MicroSnapshot | None = None, ) -> None: q = None if quote_id and quote_id in self._quotes: q = self._quotes[quote_id] elif trade_id is not None and trade_id in self._quotes_by_trade: q = self._quotes.get(self._quotes_by_trade[trade_id]) if q is None or q.status != "open": return now = time.time() q.status = "canceled" ev = { "event": "quote_canceled", "quote_id": q.quote_id, "pair": q.pair, "side": q.side, "quote_price": q.quote_price, "quote_created_time": _iso(q.created_ts), "quote_cancel_time": _iso(now), "quote_cancel_epoch": now, "time_alive_sec": now - q.created_ts, "cancel_reason": reason, "filled": False, "status": "canceled", "trade_id": q.trade_id, } if snap: ev.update(snap.to_book_fields()) self.write(ev) def bind_trade(self, quote_id: str, trade_id: int) -> None: if quote_id in self._quotes: self._quotes[quote_id].trade_id = trade_id self._quotes_by_trade[trade_id] = quote_id # ------------------------------------------------------------------ # # Fill + path # ------------------------------------------------------------------ # def log_fill( self, pair: str, side: str, fill_price: float, amount: float, inventory: float, snap: MicroSnapshot, order_type: str = "limit", quote_id: str | None = None, trade_id: Optional[int] = None, fill_reason: str = "maker_hit", state: dict | None = None, extra: dict | None = None, ) -> str: now = time.time() fill_id = uuid.uuid4().hex[:16] # resolve quote lifecycle q: Optional[ActiveQuote] = None if quote_id and quote_id in self._quotes: q = self._quotes[quote_id] elif trade_id is not None and trade_id in self._quotes_by_trade: q = self._quotes.get(self._quotes_by_trade[trade_id]) time_to_fill = None quote_created_time = None quote_price = fill_price if q is not None: q.status = "filled" time_to_fill = now - q.created_ts quote_created_time = _iso(q.created_ts) quote_price = q.quote_price quote_id = q.quote_id det = self.book_deterioration(side, now=now, lookback=5.0) ev = { "event": "fill", "fill_id": fill_id, "quote_id": quote_id, "pair": pair, "side": side, "fill_price": fill_price, "quote_price": quote_price, "amount": amount, "inventory": inventory, "order_type": order_type, "fill_reason": fill_reason, "quote_created_time": quote_created_time, "quote_fill_time": _iso(now), "time_to_fill": time_to_fill, "trade_id": trade_id, "filled": True, } ev.update(snap.to_book_fields()) ev.update(det) if state: ev.update(state) if extra: ev.update(extra) self.write(ev) # also emit quote_filled lifecycle event if q is not None: self.write( { "event": "quote_filled", "quote_id": q.quote_id, "fill_id": fill_id, "pair": pair, "side": q.side, "quote_price": q.quote_price, "quote_created_time": _iso(q.created_ts), "quote_fill_time": _iso(now), "time_to_fill": time_to_fill, "fill_reason": fill_reason, "filled": True, "status": "filled", "trade_id": trade_id, **snap.to_book_fields(), **det, } ) self._pending[fill_id] = PendingFillPath( fill_id=fill_id, pair=pair, side=side, fill_price=fill_price, fill_ts=now, quote_id=quote_id, ) return fill_id def attach_exit_reason(self, fill_id: str, exit_reason: str) -> None: if fill_id in self._pending: self._pending[fill_id].exit_reason = exit_reason # also write lightweight annotation self.write( { "event": "fill_exit", "fill_id": fill_id, "exit_reason": exit_reason, } ) def update_paths(self, pair: str, last_price: float, now: float | None = None) -> None: now = now or time.time() finished = [] for fid, p in self._pending.items(): if p.pair != pair or p.done: continue p.min_price = min(p.min_price, last_price) p.max_price = max(p.max_price, last_price) mae, mfe = p.signed_excursions() age = now - p.fill_ts def mark(horizon_attr_price, horizon_mae, horizon_mfe, sec, price_val): if getattr(p, horizon_attr_price) is None and age >= sec: setattr(p, horizon_attr_price, price_val) setattr(p, horizon_mae, mae) setattr(p, horizon_mfe, mfe) mark("after_10s_price", "mae_10s", "mfe_10s", 10, last_price) mark("after_30s_price", "mae_30s", "mfe_30s", 30, last_price) mark("after_1m_price", "mae_1m", "mfe_1m", 60, last_price) if p.after_5m_price is None and age >= 300: p.after_5m_price = last_price p.mae_5m = mae p.mfe_5m = mfe p.done = True # Price MAE absolute if p.side == "long": price_mae = p.min_price - p.fill_price price_mfe = p.max_price - p.fill_price else: price_mae = p.fill_price - p.max_price # negative if adverse up price_mfe = p.fill_price - p.min_price self.write( { "event": "fill_path", "fill_id": p.fill_id, "quote_id": p.quote_id, "pair": p.pair, "side": p.side, "fill_price": p.fill_price, "exit_reason": p.exit_reason, "after_10s_price": p.after_10s_price, "after_30s_price": p.after_30s_price, "after_1m_price": p.after_1m_price, "after_5m_price": p.after_5m_price, "min_price": p.min_price, "max_price": p.max_price, # percent "mae_10s": p.mae_10s, "mae_30s": p.mae_30s, "mae_1m": p.mae_1m, "mae_5m": p.mae_5m, "mfe_10s": p.mfe_10s, "mfe_30s": p.mfe_30s, "mfe_1m": p.mfe_1m, "mfe_5m": p.mfe_5m, # absolute price "price_mae": price_mae, "price_mfe": price_mfe, "price_mae_pct": mae, "price_mfe_pct": mfe, } ) finished.append(fid) for fid in finished: self._pending.pop(fid, None) @property def pending_count(self) -> int: return len(self._pending) # 兼容旧 API def log_quote(self, *args, **kwargs): """Deprecated wrapper → create_quote for live quotes; heartbeat uses book only.""" return self.create_quote(*args, **kwargs)