From 5bc6e3e18534ea3bd319693a3746cca86905666c Mon Sep 17 00:00:00 2001 From: Umais Khan Date: Sat, 4 Apr 2026 17:11:31 +0500 Subject: [PATCH] Create market_data.py --- core/market_data.py | 305 ++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 305 insertions(+) create mode 100644 core/market_data.py diff --git a/core/market_data.py b/core/market_data.py new file mode 100644 index 0000000..b2d1196 --- /dev/null +++ b/core/market_data.py @@ -0,0 +1,305 @@ +import asyncio +import logging +import time +from collections import deque +from datetime import datetime, timedelta, timezone +from typing import Any, Dict, List, Optional, Tuple + +from core.state import StateManager +from brokers.mt5_controller import MT5Controller +from core.logger import get_logger +try: + from analysis.market_microstructure import AdvancedMarketMicrostructure, TickData # type: ignore +except Exception: + AdvancedMarketMicrostructure = None # type: ignore + TickData = None # type: ignore +try: + from analysis.microstructure_analyzer import AdvancedMicrostructureAnalyzer # type: ignore +except Exception: + AdvancedMicrostructureAnalyzer = None # type: ignore +try: + from core.hft_optimizer import HFTOptimizer # type: ignore +except Exception: + HFTOptimizer = None # type: ignore + + +class _OHLCAggregator: + def __init__(self, interval_seconds: int) -> None: + self.interval = interval_seconds + self.bucket_start: Optional[int] = None + self.o: Optional[float] = None + self.h: Optional[float] = None + self.l: Optional[float] = None + self.c: Optional[float] = None + self.v: float = 0.0 + + def _bucket_for(self, ts: int) -> int: + return ts - (ts % self.interval) + + def update(self, price: float, volume: float, ts: int) -> Tuple[Optional[Dict[str, Any]], Dict[str, Any]]: + b = self._bucket_for(ts) + finished: Optional[Dict[str, Any]] = None + if self.bucket_start is None: + self.bucket_start = b + self.o = price + self.h = price + self.l = price + self.c = price + self.v = volume + elif b != self.bucket_start: + finished = { + "t": self.bucket_start, + "open": self.o, + "high": self.h, + "low": self.l, + "close": self.c, + "volume": self.v, + } + self.bucket_start = b + self.o = price + self.h = price + self.l = price + self.c = price + self.v = volume + else: + if self.h is None or price > self.h: + self.h = price + if self.l is None or price < self.l: + self.l = price + self.c = price + self.v += volume + current = { + "t": self.bucket_start, + "open": self.o, + "high": self.h, + "low": self.l, + "close": self.c, + "volume": self.v, + } + return finished, current + + +class _MultiTFBuilder: + def __init__(self, tf_map: Dict[str, int]) -> None: + self._tf_map = tf_map + self._aggs: Dict[str, _OHLCAggregator] = {tf: _OHLCAggregator(sec) for tf, sec in tf_map.items()} + + def update(self, price: float, volume: float, ts: int) -> Tuple[List[Tuple[str, Dict[str, Any]]], Dict[str, Dict[str, Any]]]: + finished: List[Tuple[str, Dict[str, Any]]] = [] + currents: Dict[str, Dict[str, Any]] = {} + for tf, agg in self._aggs.items(): + f, c = agg.update(price, volume, ts) + if f: + finished.append((tf, f)) + currents[tf] = c + return finished, currents + + +class MarketDataManager: + def __init__(self, state: StateManager, controller: MT5Controller) -> None: + self.state = state + self.controller = controller + self._loop = asyncio.get_event_loop() + self._tf_seconds: Dict[str, int] = {"M1": 60, "M5": 300, "M15": 900, "M30": 1800, "H1": 3600} + self._builder = _MultiTFBuilder(self._tf_seconds) + self.symbol_data: Dict[str, Dict[str, deque]] = { + self.controller.symbol: {tf: deque(maxlen=500) for tf in self._tf_seconds.keys()} + } + self._logger = get_logger("market_data") + try: + self._ms = AdvancedMarketMicrostructure(self.controller.symbol, state=self.state) if AdvancedMarketMicrostructure else None + except Exception: + self._ms = None + try: + self._msa = AdvancedMicrostructureAnalyzer(self.controller.symbol, state=self.state) if AdvancedMicrostructureAnalyzer else None + except Exception: + self._msa = None + try: + self._hft = HFTOptimizer(self.controller.symbol, state=self.state) if HFTOptimizer else None + except Exception: + self._hft = None + self._hft_step_ctr: int = 0 + + def _update_all_exec(self, price: float, volume: float, ts: int) -> Tuple[List[Tuple[str, Dict[str, Any]]], Dict[str, Dict[str, Any]]]: + return self._builder.update(price, volume, ts) + + def _validate_candle(self, tf: str, c: Dict[str, Any], last_t: Optional[int]) -> bool: + try: + o = float(c["open"]) + h = float(c["high"]) + l = float(c["low"]) + cl = float(c["close"]) + t = int(c["t"]) + except Exception: + self._logger.warning("invalid candle values for %s: %s", tf, c) + return False + if not (h >= l and h >= max(o, cl) and l <= min(o, cl)): + self._logger.warning("price validity failed for %s: %s", tf, c) + return False + if last_t is not None and t <= last_t: + self._logger.warning("timestamp order/duplicate failed for %s: last=%s new=%s", tf, last_t, t) + return False + sec = self._tf_seconds.get(tf, 0) + if sec and (t % sec) != 0: + self._logger.warning("time alignment failed for %s: %s", tf, t) + return False + return True + + async def _store_candle(self, tf: str, candle: Dict[str, Any]) -> None: + dq = self.symbol_data[self.controller.symbol][tf] + last_t = dq[-1]["t"] if len(dq) > 0 else None + if self._validate_candle(tf, candle, last_t): + dq.append(candle) + await self.state.push_candle_update_for(self.controller.symbol, tf, candle) + else: + self._logger.warning("skipped corrupted candle for %s", tf) + + async def load_historical_data(self, bars: int = 500) -> None: + mt5 = self.controller._import_mt5() # type: ignore[attr-defined] + if mt5 is None: + self._logger.error("mt5 not available for historical load") + return + end = datetime.now(timezone.utc) + for tf, sec in self._tf_seconds.items(): + try: + start = end - timedelta(seconds=sec * bars) + tf_const_name = "TIMEFRAME_" + tf + tf_const = getattr(mt5, tf_const_name) + def _copy(): + return mt5.copy_rates_range(self.controller.symbol, tf_const, start, end) + rates = await self._loop.run_in_executor(None, _copy) + if rates is None: + self._logger.warning("no rates for %s", tf) + continue + # Ensure deques are empty for fresh load + dq = self.symbol_data[self.controller.symbol][tf] + dq.clear() + last_t: Optional[int] = None + names = set(getattr(rates, "dtype", ()).names or ()) + for r in rates: + rt = int(r["time"]) + boundary = rt - (rt % sec) + if "real_volume" in names: + vol_val = float(r["real_volume"]) + else: + vol_val = float(r["tick_volume"]) + c = { + "t": boundary, + "open": float(r["open"]), + "high": float(r["high"]), + "low": float(r["low"]), + "close": float(r["close"]), + "volume": vol_val, + } + if self._validate_candle(tf, c, last_t): + dq.append(c) + last_t = c["t"] + else: + self._logger.warning("historical candle rejected %s: %s", tf, c) + seed = list(dq)[-min(len(dq), 300 if tf == "M1" else 200):] + for c in seed: + await self.state.push_candle_update_for(self.controller.symbol, tf, c) + except Exception as e: + self._logger.error("historical load error %s for %s", e, tf) + + async def run(self) -> None: + # Start trade listener to feed microstructure analyzers + try: + tq = await self.state.add_trade_listener() + async def _forward_trades(): + while True: + try: + msg = await tq.get() + tr = (msg or {}).get("trade") or {} + sym = str((msg or {}).get("symbol") or tr.get("symbol") or self.controller.symbol) + if sym != self.controller.symbol: + continue + try: + if self._ms: + self._ms.ingest_trade(tr) + except Exception: + pass + try: + if self._msa: + self._msa.ingest_trade(tr) + except Exception: + pass + try: + if getattr(self, "_hft", None): + sz = float((tr or {}).get("lot") or (tr or {}).get("size") or 0.0) + sd = str((tr or {}).get("signal") or (tr or {}).get("side") or "") + self._hft.liqdet.update(sz, sd, int((tr or {}).get("timestamp") or int(time.time()))) + except Exception: + pass + except asyncio.CancelledError: + break + except Exception: + await asyncio.sleep(0.05) + asyncio.create_task(_forward_trades()) + except Exception: + pass + async for tick in self.controller.tick_stream(): + try: + bid = tick.get("bid") + ask = tick.get("ask") + last = tick.get("last") or bid or ask + volume = float(tick.get("volume") or 0.0) + ts = int(tick.get("timestamp")) + price = float(last) + await self.state.set_latest_tick_for(self.controller.symbol, tick) + try: + if self._ms and TickData: + td = TickData(symbol=self.controller.symbol, timestamp=ts, bid=float(bid or 0.0), ask=float(ask or 0.0), last=float(last or 0.0), volume=float(volume or 0.0)) + self._ms.ingest_tick(td) + else: + pass + except Exception: + pass + try: + if self._msa: + self._msa.ingest_tick({"symbol": self.controller.symbol, "timestamp": ts, "bid": float(bid or 0.0), "ask": float(ask or 0.0), "last": float(last or 0.0), "volume": float(volume or 0.0)}) + else: + pass + except Exception: + pass + try: + if self._hft: + self._hft.update_tick(float(bid or 0.0), float(ask or 0.0), float(last or 0.0), ts=int(ts), volume=float(volume or 0.0)) + self._hft_step_ctr += 1 + if (self._hft_step_ctr % 5) == 0: + self._hft.step() + else: + pass + except Exception: + pass + finished, currents = await self._loop.run_in_executor(None, lambda: self._update_all_exec(price, volume, ts)) + for tf, c in finished: + await self._store_candle(tf, c) + h1 = currents.get("H1") + m5 = currents.get("M5") + if h1: + await self.state.set_latest_ohlc_1h(h1) + if m5: + await self.state.set_latest_ohlc_5m(m5) + self._logger.debug("tick %s %s", tick.get("symbol"), price) + except asyncio.CancelledError: + break + except Exception as e: + self._logger.error("market data error %s", e) + + async def account_monitor(self) -> None: + while True: + try: + info = await self.controller.get_account_info() + if info is not None: + await self.state.set_account_info_for(self.controller.symbol, info) + try: + self._logger.info("account balance %s equity %s", info.get("balance"), info.get("equity")) + except Exception: + pass + await asyncio.sleep(2.0) + except asyncio.CancelledError: + break + except Exception as e: + self._logger.error("account monitor error %s", e) + await asyncio.sleep(2.0)