From ab7c59d387a0150bd716fac34d0075e898b1a1e1 Mon Sep 17 00:00:00 2001 From: Umais Khan Date: Sat, 4 Apr 2026 17:09:54 +0500 Subject: [PATCH] Create volume_profile.py --- core/volume_profile.py | 355 +++++++++++++++++++++++++++++++++++++++++ 1 file changed, 355 insertions(+) create mode 100644 core/volume_profile.py diff --git a/core/volume_profile.py b/core/volume_profile.py new file mode 100644 index 0000000..32d678d --- /dev/null +++ b/core/volume_profile.py @@ -0,0 +1,355 @@ +import asyncio +import json +import os +import time +from typing import Any, Dict, List, Optional, Tuple, Union + +import numpy as np + +from core.logger import get_logger +from core.state import StateManager + + +def _group_ranges(indices: List[int]) -> List[Tuple[int, int]]: + if not indices: + return [] + indices.sort() + ranges: List[Tuple[int, int]] = [] + start = indices[0] + prev = indices[0] + for i in indices[1:]: + if i == prev + 1: + prev = i + continue + ranges.append((start, prev)) + start = i + prev = i + ranges.append((start, prev)) + return ranges + + +class VolumeProfile: + def __init__(self, state: Optional[StateManager] = None, symbol: str = "EURUSDm", bins: int = 50, window: int = 500, enabled: bool = True, overlay_enabled: bool = True, dashboard_enabled: bool = True) -> None: + self.state = state + self.symbol = symbol + self._bins = int(bins) + self._window = int(window) + self._enabled = bool(enabled) + self._overlay_enabled = bool(overlay_enabled) + self._dashboard_enabled = bool(dashboard_enabled) + self._logger = get_logger("volume_profile") + self._loop = asyncio.get_event_loop() + self._last_profile: Optional[Dict[str, Any]] = None + self._last_levels: Optional[Dict[str, Any]] = None + self._config_loaded = False + self._bin_width: float = 0.0 + self._timeframe = "M5" + self._load_config() + + def _load_config(self) -> None: + try: + path = os.path.join("config", "runtime_config.json") + if os.path.exists(path): + with open(path, "r", encoding="utf-8") as f: + raw = json.load(f) + ev = raw.get("enablevolumeprofile") + if ev is not None: + self._enabled = bool(ev) + vb = raw.get("vpbins") + if vb is not None: + self._bins = int(vb) + vw = raw.get("vpwindow") + if vw is not None: + self._window = int(vw) + vo = raw.get("vpoverlay") + if vo is not None: + self._overlay_enabled = bool(vo) + self._config_loaded = True + except Exception: + self._config_loaded = False + + def compute(self, data: List[Dict[str, Any]], bins: int = 50, window: int = 500) -> Dict[str, Any]: + if not data: + return {"edges": np.array([]), "volumes": np.array([]), "centers": np.array([]), "normalized": np.array([])} + w = max(1, int(window)) + d = data[-w:] + lows = np.array([float(x.get("low") or x.get("close") or 0.0) for x in d], dtype=float) + highs = np.array([float(x.get("high") or x.get("close") or 0.0) for x in d], dtype=float) + closes = np.array([float(x.get("close") or 0.0) for x in d], dtype=float) + vols = np.array([float(x.get("volume") or 0.0) for x in d], dtype=float) + tpx = (highs + lows + closes) / 3.0 + lo = float(np.min(lows)) if lows.size else float(np.min(closes) if closes.size else 0.0) + hi = float(np.max(highs)) if highs.size else float(np.max(closes) if closes.size else 0.0) + if not np.isfinite(lo) or not np.isfinite(hi) or hi <= lo: + return {"edges": np.array([]), "volumes": np.array([]), "centers": np.array([]), "normalized": np.array([])} + b = max(5, int(bins)) + hist, edges = np.histogram(tpx, bins=b, range=(lo, hi), weights=vols) + centers = (edges[:-1] + edges[1:]) / 2.0 + vmax = float(np.max(hist)) if hist.size else 0.0 + norm = (hist / (vmax or 1.0)) if hist.size else np.array([]) + self._bin_width = float((edges[1] - edges[0]) if len(edges) >= 2 else 0.0) + return {"edges": edges, "volumes": hist, "centers": centers, "normalized": norm} + + def get_poc(self, profile: Dict[str, Any]) -> Optional[float]: + vols = profile.get("volumes") + centers = profile.get("centers") + if vols is None or centers is None: + return None + if len(vols) == 0 or len(centers) == 0: + return None + i = int(np.argmax(vols)) + return float(centers[i]) + + def gethvnlvn(self, profile: Dict[str, Any]) -> Tuple[List[Tuple[float, float]], List[Tuple[float, float]]]: + vols = profile.get("volumes") + edges = profile.get("edges") + if vols is None or edges is None or len(vols) == 0 or len(edges) < 2: + return [], [] + v = np.array(vols, dtype=float) + nz = v[v > 0.0] + if nz.size == 0: + return [], [] + hv_thr = float(np.percentile(nz, 80.0)) + lv_thr = float(np.percentile(nz, 20.0)) + hv_idx = [i for i in range(len(v)) if v[i] >= hv_thr] + lv_idx = [i for i in range(len(v)) if v[i] <= lv_thr] + hv_groups = _group_ranges(hv_idx) + lv_groups = _group_ranges(lv_idx) + hvn: List[Tuple[float, float]] = [] + lvn: List[Tuple[float, float]] = [] + for a, b in hv_groups: + hvn.append((float(edges[a]), float(edges[b + 1]))) + for a, b in lv_groups: + lvn.append((float(edges[a]), float(edges[b + 1]))) + if len(hvn) > 5: + scores = [sum(vols[i:j + 1]) for i, j in hv_groups] + order = np.argsort(scores)[::-1][:5].tolist() + hvn = [hvn[i] for i in order] + if len(lvn) > 5: + scores_l = [sum(vols[i:j + 1]) for i, j in lv_groups] + order_l = np.argsort(scores_l)[:5].tolist() + lvn = [lvn[i] for i in order_l] + return hvn, lvn + + def overlay(self, frame: Optional[Any], profile: Dict[str, Any]) -> Any: + try: + import matplotlib.pyplot as plt + except Exception: + return None + edges = profile.get("edges") + norm = profile.get("normalized") + centers = profile.get("centers") + if edges is None or norm is None or len(edges) < 2 or len(norm) == 0: + return None + poc = self.get_poc(profile) + hvn, lvn = self.gethvnlvn(profile) + if frame is None: + fig, ax = plt.subplots(figsize=(4, 6)) + else: + fig, ax = (None, frame) + yvals = (edges[:-1] + edges[1:]) / 2.0 + ax.barh(yvals, norm, height=(edges[1] - edges[0]) * 0.9, color="#4a90e2", alpha=0.6) + if poc is not None: + ax.axhline(poc, color="#ff3b30", linewidth=2.0) + for lo, hi in hvn: + ax.axhspan(lo, hi, color="#34c759", alpha=0.2) + for lo, hi in lvn: + ax.axhspan(lo, hi, color="#ffcc00", alpha=0.2) + ax.set_xlabel("Norm Vol") + ax.set_ylabel("Price") + if fig is not None: + fig.tight_layout() + return fig + return ax + + async def _gate_decisions(self) -> None: + if self.state is None: + return + q = await self.state.add_decision_listener() + while True: + try: + msg = await q.get() + if not isinstance(msg, dict): + continue + if msg.get("type") != "decision_update": + continue + if msg.get("symbol") and str(msg.get("symbol")) != self.symbol: + continue + if not self._enabled: + continue + decision = await self.state.get_latest_decision_for(self.symbol) + if not isinstance(decision, dict): + continue + if decision.get("vp_checked"): + continue + if not bool(decision.get("approved")): + continue + sig = str(decision.get("signal") or "").upper() + if sig not in ("BUY", "SELL"): + continue + levels = self._last_levels + prof = self._last_profile + if not levels or not prof: + continue + poc = float(levels.get("poc") or 0.0) + hvn = levels.get("hvn") or [] + lvn = levels.get("lvn") or [] + try: + candles = await self.state.get_latest_candles_m5_for(self.symbol) + except Exception: + candles = [] + px = float((candles[-1].get("close") if candles else 0.0) or 0.0) + bw = float(self._bin_width or 0.0) + if bw <= 0.0 and prof.get("edges") is not None and len(prof.get("edges")) >= 2: + ed = prof.get("edges") + bw = float(ed[1] - ed[0]) + + # Logic encrypted for proprietary alpha execution + + ok = True + + # Logic encrypted for proprietary alpha execution + + decision["vp_checked"] = True + decision["vp_levels"] = {"poc": poc, "hvn": hvn, "lvn": lvn} + try: + await self.state.set_latest_decision_for(self.symbol, decision) + except Exception: + pass + except asyncio.CancelledError: + break + except Exception: + await asyncio.sleep(0.05) + + async def _compute_loop(self) -> None: + if self.state is None: + return + q = await self.state.add_candle_listener() + while True: + try: + msg = await q.get() + if not isinstance(msg, dict): + continue + if msg.get("type") != "candle_update": + continue + if msg.get("symbol") and str(msg.get("symbol")) != self.symbol: + continue + tf = str(msg.get("timeframe") or "") + if tf != self._timeframe: + continue + if not self._enabled: + continue + try: + candles = await self.state.get_latest_candles_m5_for(self.symbol) + except Exception: + candles = [] + def _calc(): + return self.compute(candles, bins=self._bins, window=self._window) + profile = await asyncio.get_event_loop().run_in_executor(None, _calc) + self._last_profile = profile + poc = self.get_poc(profile) + hvn, lvn = self.gethvnlvn(profile) + levels = {"poc": float(poc) if poc is not None else None, "hvn": hvn, "lvn": lvn} + self._last_levels = levels + ts = int(msg.get("time") or int(time.time())) + try: + self._logger.info("[VP] Computed profile with %d bins", int(len(profile.get("volumes") or []))) + except Exception: + pass + try: + self._logger.info("[VP] POC=%s, HVN=%s, LVN=%s", str(levels.get("poc")), str(hvn), str(lvn)) + except Exception: + pass + try: + payload = {"timestamp": ts, "bins": self._bins, "window": self._window, "poc": levels.get("poc"), "hvn": hvn, "lvn": lvn} + self._logger.info(json.dumps(payload)) + except Exception: + pass + try: + await self.state.set_latest_risk_features_for(self.symbol, {"vp_poc": levels.get("poc"), "vp_hvn": hvn, "vp_lvn": lvn}) + except Exception: + pass + if self._dashboard_enabled: + try: + ed = profile.get("edges") or [] + norm = profile.get("normalized") or [] + centers = profile.get("centers") or [] + pts = [] + try: + c = centers.tolist() if hasattr(centers, "tolist") else list(centers) + n = norm.tolist() if hasattr(norm, "tolist") else list(norm) + for i in range(min(len(c), len(n))): + pts.append([float(c[i]), float(n[i])]) + except Exception: + pts = [] + payload_dbg = {"stage": "volume_profile", "timestamp": ts, "bins": self._bins, "window": self._window, "poc": levels.get("poc"), "hvn": hvn, "lvn": lvn, "normalized_points": pts, "bin_width": float(self._bin_width or 0.0)} + await self.state.set_latest_pipeline_debug_for(self.symbol, payload_dbg) + except Exception: + pass + except asyncio.CancelledError: + break + except Exception: + await asyncio.sleep(0.05) + + async def _config_loop(self) -> None: + if self.state is None: + return + q = await self.state.add_config_listener() + while True: + try: + msg = await q.get() + if not isinstance(msg, dict): + continue + if msg.get("type") != "config_update": + continue + p = str(msg.get("parameter") or "").lower() + v = msg.get("value") + if p in ("enablevolumeprofile", "volumeprofile_enabled", "vp_enabled"): + try: + self._enabled = bool(v if isinstance(v, bool) else (float(v) > 0.5 if isinstance(v, (int, float)) else str(v).lower() in ("1", "true", "yes", "on"))) + except Exception: + pass + elif p in ("vpbins", "vp_bins"): + try: + self._bins = max(5, int(float(v))) + except Exception: + pass + elif p in ("vpwindow", "vp_window"): + try: + self._window = max(50, int(float(v))) + except Exception: + pass + elif p in ("vpoverlay", "vp_overlay"): + try: + self._overlay_enabled = bool(v if isinstance(v, bool) else (float(v) > 0.5 if isinstance(v, (int, float)) else str(v).lower() in ("1", "true", "yes", "on"))) + except Exception: + pass + except asyncio.CancelledError: + break + except Exception: + await asyncio.sleep(0.05) + + async def run(self) -> None: + if self.state is None: + return + if not self._config_loaded: + self._load_config() + tasks = [ + asyncio.create_task(self._compute_loop()), + asyncio.create_task(self._gate_decisions()), + asyncio.create_task(self._config_loop()), + ] + try: + await asyncio.gather(*tasks) + finally: + for t in tasks: + try: + t.cancel() + except Exception: + pass + + def latest_profile(self) -> Optional[Dict[str, Any]]: + return self._last_profile + + def latest_levels(self) -> Optional[Dict[str, Any]]: + return self._last_levels