Files
2026-04-04 17:09:54 +05:00

356 lines
15 KiB
Python

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