Files
2026-04-04 17:11:31 +05:00

306 lines
13 KiB
Python

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)