Spaces:
Sleeping
Sleeping
| # app.py — ULTRA‑X Multi‑Source Signal Board Pro (Mega Edition) | |
| # Fusion v4.0 — dynamic weights, regime detection, Pro+, Mega, DataMonster | |
| import os | |
| import re | |
| import json | |
| import time | |
| from typing import Any, Dict, List, Optional, Tuple | |
| from collections import deque, defaultdict | |
| from datetime import datetime | |
| from zoneinfo import ZoneInfo | |
| import numpy as np | |
| import pandas as pd | |
| import gradio as gr | |
| from fastapi import FastAPI, HTTPException, Query, Request, WebSocket, WebSocketDisconnect | |
| from fastapi.middleware.cors import CORSMiddleware | |
| from fastapi.responses import StreamingResponse | |
| # ========================= | |
| # ENV / CONFIG | |
| # ========================= | |
| WEBHOOK_TOKEN = os.getenv("WEBHOOK_TOKEN", "") | |
| BODY_SECRETS = {s.strip() for s in os.getenv("BODY_SECRETS", "").split(",") if s.strip()} | |
| MAX_ALERTS = int(os.getenv("MAX_ALERTS", "20000")) | |
| PERSIST_PATH = os.getenv("PERSIST_PATH", "") | |
| MAX_LOG_MB = int(os.getenv("MAX_LOG_MB", "50")) | |
| ATHENS_TZ = ZoneInfo("Europe/Athens") | |
| RUNTIME_CONFIG: Dict[str, Any] = { | |
| "fusion": { | |
| "time_window_minutes": 5, | |
| "min_confidence": 0.60, | |
| "weight_threshold": 0.15, | |
| "consensus_threshold": 0.70, | |
| }, | |
| "weights": { | |
| "ULTRA-X-PRECISION90": 0.95, | |
| "MASTER_BTC_PRO": 0.85, | |
| "QSC_PRO": 0.80, | |
| "ULTRAX_FEEDER": 0.75, | |
| "SQC_PRO": 0.70, | |
| "RTE": 0.65, | |
| "BOS_FORMED": 0.60 | |
| } | |
| } | |
| ALERT_PATTERNS = { | |
| "ULTRAX_FEEDER": [ | |
| r"ULTRAX\s*Feeder", | |
| r"ULTRAX\s*Feeder\s*\(ST.*%R.*AI-proxy\)" | |
| ], | |
| "ULTRA-X-PRECISION90": [ | |
| r"PTP\s*Data\s*Emitter\s*ULTRA-X\s*Precision90", | |
| r"ULTRA-X\s*Precision\s*V?9?" | |
| ], | |
| "MASTER_BTC_PRO": [ | |
| r"Master\s*BTC\s*Indicator\s*Pro\s*Enhanced", | |
| r"Master\s*BTC\s*Indicator\s*Pro", | |
| r"Master\s*BTC\b" | |
| ], | |
| "QSC_PRO": [r"\bQSC[_ ]?Pro\b", r"\bQSC_PRO\b"], | |
| "SQC_PRO": [r"\bSQC[_ ]?Pro\b", r"\bSQC_PRO\b"], | |
| "RTE": [r"%R\s*Trend\s*Exhaustion", r"\bRTE\b"], | |
| "BOS_FORMED": [r"Internal\s*(Bullish|Bearish)\s*BOS\s*formed", r"\bBOS\s*formed\b"] | |
| } | |
| TIMEFRAME_MAP = { | |
| "1m": "1", "3m": "3", "5m": "5", "15m": "15", "30m": "30", "45m": "45", | |
| "1h": "60", "2h": "120", "4h": "240", "6h": "360", "12h": "720", | |
| "1d": "1D", "1w": "1W" | |
| } | |
| LONG_WORDS = {"LONG", "BUY", "BULL", "BULLISH", "UP"} | |
| SHORT_WORDS = {"SHORT", "SELL", "BEAR", "BEARISH", "DOWN"} | |
| # ========================= | |
| # HELPERS | |
| # ========================= | |
| def _now_ms() -> int: | |
| return int(time.time() * 1000) | |
| def _to_float(x: Any, default: Optional[float] = None) -> Optional[float]: | |
| try: | |
| if x is None: | |
| return default | |
| return float(x) | |
| except Exception: | |
| return default | |
| def ms_to_iso_utc(ms: Any) -> Optional[str]: | |
| try: | |
| return time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime(int(ms) / 1000)) | |
| except Exception: | |
| return None | |
| def iso_to_ms(s: str) -> Optional[int]: | |
| try: | |
| if s.endswith("Z"): | |
| s = s.replace("Z", "+00:00") | |
| return int(datetime.fromisoformat(s).timestamp() * 1000) | |
| except Exception: | |
| return None | |
| def athens_str_from_ms(ms: int) -> str: | |
| try: | |
| dt = datetime.fromtimestamp(ms / 1000, tz=ATHENS_TZ) | |
| return dt.strftime("%Y-%m-%d %H:%M:%S") | |
| except Exception: | |
| return "" | |
| def persist_append(row: Dict[str, Any]) -> None: | |
| if not PERSIST_PATH: | |
| return | |
| try: | |
| os.makedirs(os.path.dirname(PERSIST_PATH), exist_ok=True) | |
| if os.path.exists(PERSIST_PATH) and (os.path.getsize(PERSIST_PATH) / (1024 * 1024)) > MAX_LOG_MB: | |
| ts = time.strftime("%Y%m%d_%H%M%S") | |
| os.rename(PERSIST_PATH, f"{PERSIST_PATH}.{ts}.bak") | |
| with open(PERSIST_PATH, "a", encoding="utf-8") as f: | |
| f.write(json.dumps(row, ensure_ascii=False) + "\n") | |
| except Exception: | |
| pass | |
| def normalize_alert_name(alert_name: Optional[str]) -> Optional[str]: | |
| if not alert_name: | |
| return None | |
| for std, pats in ALERT_PATTERNS.items(): | |
| for pat in pats: | |
| if re.search(pat, alert_name, flags=re.IGNORECASE): | |
| return std | |
| cleaned = re.sub(r"\([^)]*\)", "", alert_name).strip() | |
| cleaned = re.sub(r"[^a-zA-Z0-9_]+", "_", cleaned).upper() | |
| return cleaned or None | |
| def extract_ticker_tf_from_pair(pair_tf: Optional[str]) -> Tuple[Optional[str], Optional[str]]: | |
| if not pair_tf: | |
| return None, None | |
| parts = [p.strip() for p in str(pair_tf).split(",")] | |
| ticker = parts[0] if parts else None | |
| tf = None | |
| if len(parts) > 1: | |
| key = parts[1].lower() | |
| tf = TIMEFRAME_MAP.get(key, key) | |
| return ticker, tf | |
| def parse_any_time(val: Any) -> Tuple[Optional[int], Optional[str]]: | |
| if val is None: | |
| return None, None | |
| s = str(val).strip() | |
| if s.isdigit(): | |
| n = int(s) | |
| if n < 10_000_000_000: | |
| n *= 1000 | |
| return n, ms_to_iso_utc(n) | |
| try: | |
| if "T" in s and "Z" in s: | |
| ms = iso_to_ms(s) | |
| return ms, s | |
| except Exception: | |
| pass | |
| try: | |
| if "'" in s: | |
| s2 = s.replace("'", "20") | |
| for fmt in ("%a %d %b %Y %H:%M:%S", "%a %d %b %y %H:%M:%S"): | |
| try: | |
| dt = datetime.strptime(s2, fmt) | |
| ms = int(dt.timestamp() * 1000) | |
| return ms, ms_to_iso_utc(ms) | |
| except ValueError: | |
| continue | |
| except Exception: | |
| pass | |
| return None, None | |
| def text_to_signal(txt: str) -> str: | |
| s = (txt or "").upper() | |
| if any(w in s for w in LONG_WORDS): | |
| return "LONG" | |
| if any(w in s for w in SHORT_WORDS): | |
| return "SHORT" | |
| return "FLAT" | |
| # ========================= | |
| # ULTRA‑X tech fallback | |
| # ========================= | |
| def _strength_from_payload(data: Dict[str, Any], sig_type: str) -> int: | |
| strength = 0 | |
| trend_score = abs(int(data.get("trend_score", 0))) | |
| strength += min(trend_score * 10, 30) | |
| adx_val = _to_float(data.get("adx_val"), 0) or 0 | |
| strength += min(adx_val, 30) | |
| vol_ratio = _to_float(data.get("vol_ratio"), 0) or 0 | |
| if vol_ratio > 2.0: | |
| strength += 20 | |
| elif vol_ratio > 1.5: | |
| strength += 10 | |
| rsi_val = _to_float(data.get("rsi"), 50) or 50 | |
| if sig_type == "long" and rsi_val < 65: | |
| strength += 10 | |
| if sig_type == "short" and rsi_val > 35: | |
| strength += 10 | |
| return int(min(strength, 100)) | |
| def compute_signal_ultra(payload: Dict[str, Any]) -> Dict[str, Any]: | |
| long_c = bool(payload.get("long_condition", False)) | |
| short_c = bool(payload.get("short_condition", False)) | |
| if long_c and not short_c: | |
| return {"signal": "LONG", "strength": _strength_from_payload(payload, "long")} | |
| if short_c and not long_c: | |
| return {"signal": "SHORT", "strength": _strength_from_payload(payload, "short")} | |
| return {"signal": "FLAT", "strength": 0} | |
| # ========================= | |
| # SCHEMA / CONFIDENCE | |
| # ========================= | |
| def detect_schema(p: Dict[str, Any]) -> str: | |
| if isinstance(p.get("signal"), str): | |
| return "GENERIC_SIG" | |
| if isinstance(p.get("bullish"), bool) or isinstance(p.get("bearish"), bool): | |
| return "GENERIC_SIG" | |
| src_upper = str(p.get("src") or p.get("emitter") or p.get("source") or p.get("indicator") or "").upper() | |
| if p.get("source") == "tradingview" or "ULTRA" in src_upper or "EMITTER" in src_upper: | |
| return "ULTRA_X" | |
| if src_upper in {"SQC_PRO", "QSC_PRO", "MASTER_BTC_PRO", "ULTRAX_FEEDER", "RTE", "BOS_FORMED"}: | |
| return "GENERIC_SIG" | |
| return "UNKNOWN" | |
| def extract_confidence(p: Dict[str, Any]) -> Optional[float]: | |
| for k in ["confidence", "strength", "score", "prob", "probability", "conf", "weight"]: | |
| if k in p: | |
| try: | |
| return float(p[k]) | |
| except Exception: | |
| continue | |
| if p.get("bullish") is True or p.get("bearish") is True: | |
| return 1.0 | |
| return None | |
| def calc_rr(price: Optional[float], sl: Optional[float], tp: Optional[float], signal: str) -> Optional[float]: | |
| try: | |
| if price is None or sl is None or tp is None: | |
| return None | |
| signal = (signal or "").upper() | |
| if signal.startswith("LONG"): | |
| risk, reward = price - sl, tp - price | |
| elif signal.startswith("SHORT"): | |
| risk, reward = sl - price, price - tp | |
| else: | |
| return None | |
| return round(reward / risk, 2) if risk and risk > 0 else None | |
| except Exception: | |
| return None | |
| # ========================= | |
| # PERFORMANCE (dynamic weights) | |
| # ========================= | |
| class SourcePerformanceTracker: | |
| def __init__(self): | |
| self.source_scores = defaultdict(lambda: {"total_signals": 0, "successful_signals": 0}) | |
| def _evaluate_signal_success(signal: str, actual_price_move: float) -> bool: | |
| s = (signal or "").upper() | |
| if s.startswith("LONG"): | |
| return actual_price_move > 0 | |
| if s.startswith("SHORT"): | |
| return actual_price_move < 0 | |
| return False | |
| def update_performance(self, source: str, signal: str, actual_price_move: float): | |
| src = source or "unknown" | |
| ok = self._evaluate_signal_success(signal, actual_price_move) | |
| self.source_scores[src]["total_signals"] += 1 | |
| if ok: | |
| self.source_scores[src]["successful_signals"] += 1 | |
| def get_source_weight(self, source: str) -> float: | |
| base_weight = float(RUNTIME_CONFIG["weights"].get(source, 0.5)) | |
| perf = self.source_scores[source] | |
| n = perf["total_signals"] | |
| if n < 10: | |
| return base_weight | |
| success_rate = perf["successful_signals"] / max(n, 1) | |
| if success_rate > 0.6: | |
| return min(base_weight * 1.3, 0.95) | |
| if success_rate < 0.4: | |
| return max(base_weight * 0.7, 0.10) | |
| return base_weight | |
| def overall_metrics(self) -> Dict[str, Any]: | |
| total = sum(v["total_signals"] for v in self.source_scores.values()) | |
| won = sum(v["successful_signals"] for v in self.source_scores.values()) | |
| win_rate = (won / total * 100) if total else 0.0 | |
| return {"total_signals": total, "wins": won, "win_rate": round(win_rate, 2)} | |
| def as_dataframe(self) -> pd.DataFrame: | |
| rows = [] | |
| for src, v in self.source_scores.items(): | |
| n = v["total_signals"] | |
| wr = (v["successful_signals"] / n * 100) if n else 0.0 | |
| rows.append({ | |
| "source": src, | |
| "total_signals": n, | |
| "successful_signals": v["successful_signals"], | |
| "win_rate_%": round(wr, 2), | |
| "current_weight": round(self.get_source_weight(src), 3) | |
| }) | |
| rows.sort(key=lambda x: (x["total_signals"], x["win_rate_%"]), reverse=True) | |
| return pd.DataFrame(rows) | |
| source_perf = SourcePerformanceTracker() | |
| def weight_strength(base_strength: int, emitter: str, confidence: Optional[float]) -> int: | |
| w = source_perf.get_source_weight(emitter) | |
| extra = 0 | |
| if confidence is not None: | |
| extra += min(int(confidence * 40), 40) if confidence <= 1 else min(int(confidence), 40) | |
| return int(min(max(base_strength + int(w * 30) + extra, 0), 100)) | |
| # ========================= | |
| # NORMALIZATION | |
| # ========================= | |
| def parse_time_fields(p: Dict[str, Any]) -> Tuple[Optional[int], Optional[str]]: | |
| for key in ["bar_time", "time", "timestamp", "ts", "t"]: | |
| if key in p and p[key] not in (None, ""): | |
| val = p[key] | |
| if isinstance(val, (int, float)) or (isinstance(val, str) and str(val).isdigit()): | |
| ms = int(float(val)) | |
| if ms < 10_000_000_000: | |
| ms *= 1000 | |
| return ms, ms_to_iso_utc(ms) | |
| if isinstance(val, str): | |
| ms = iso_to_ms(val) | |
| return ms, val | |
| return None, None | |
| def normalize_payload(p: Dict[str, Any]) -> Dict[str, Any]: | |
| schema = detect_schema(p) | |
| emitter = p.get("emitter") or p.get("src") or p.get("source") or p.get("indicator") | |
| if not emitter and p.get("alert_name"): | |
| emitter = normalize_alert_name(p.get("alert_name")) | |
| emitter = emitter or "unknown" | |
| ticker = p.get("ticker") or p.get("symbol") or p.get("s") or "" | |
| tf = str(p.get("tf") or p.get("timeframe") or p.get("interval") or p.get("resolution") or "") | |
| if (not ticker or not tf) and p.get("pair_timeframe"): | |
| t2, tf2 = extract_ticker_tf_from_pair(p.get("pair_timeframe")) | |
| ticker = ticker or (t2 or "") | |
| tf = tf or (tf2 or "") | |
| ms, iso = parse_time_fields(p) | |
| if ms is None and iso is None: | |
| ms, iso = parse_any_time(p.get("alert_time")) | |
| price = _to_float(p.get("price") if "price" in p else p.get("close") or p.get("c")) | |
| sl = _to_float(p.get("sl_level") or p.get("sl")) | |
| tp = _to_float(p.get("tp_level") or p.get("tp")) | |
| if schema == "GENERIC_SIG": | |
| if isinstance(p.get("signal"), str): | |
| signal = text_to_signal(p["signal"]) | |
| else: | |
| if p.get("bullish") is True: | |
| signal = "LONG" | |
| elif p.get("bearish") is True: | |
| signal = "SHORT" | |
| else: | |
| signal = "FLAT" | |
| conf = extract_confidence(p) or 0.0 | |
| base_strength = int(round(conf * 100)) if conf <= 1 else int(round(conf)) | |
| else: | |
| res = compute_signal_ultra(p) | |
| signal, base_strength = res["signal"], res["strength"] | |
| conf_num = extract_confidence(p) | |
| strength = weight_strength(base_strength, emitter, conf_num) | |
| rr = calc_rr(price, sl, tp, signal) | |
| return { | |
| "schema": schema, | |
| "emitter": emitter, | |
| "ticker": ticker, | |
| "tf": tf, | |
| "bar_time_ms": ms, | |
| "bar_time": iso, | |
| "price": price, | |
| "signal": signal, | |
| "strength": int(strength), | |
| "adx": _to_float(p.get("adx_val")), | |
| "rsi": _to_float(p.get("rsi")), | |
| "position_size": _to_float(p.get("position_size")), | |
| "sl": sl, | |
| "tp": tp, | |
| "rr_ratio": rr, | |
| "raw": p, | |
| } | |
| # ========================= | |
| # DATA MONSTER INGESTOR | |
| # ========================= | |
| class DataMonsterIngestor: | |
| def __init__(self): | |
| self.raw_data_buffer = deque(maxlen=100_000) | |
| self.enriched_signals = deque(maxlen=50_000) | |
| self.market_context: Dict[str, Any] = {} | |
| def ingest_from_normalized(self, norm: Dict[str, Any]) -> Dict[str, Any]: | |
| payload = dict(norm.get("raw", {})) | |
| payload.setdefault("ticker", norm.get("ticker")) | |
| payload.setdefault("tf", norm.get("tf")) | |
| payload.setdefault("emitter", norm.get("emitter")) | |
| payload.setdefault("price", norm.get("price")) | |
| payload.setdefault("signal", norm.get("signal")) | |
| payload.setdefault("strength", norm.get("strength")) | |
| ts_ms = norm.get("bar_time_ms") or _now_ms() | |
| rec = { | |
| "timestamp_ms": ts_ms, | |
| "source": norm.get("emitter"), | |
| "ticker": norm.get("ticker"), | |
| "timeframe": norm.get("tf"), | |
| "price": norm.get("price"), | |
| "signal": norm.get("signal"), | |
| "strength": norm.get("strength"), | |
| "payload": payload, | |
| } | |
| self.raw_data_buffer.append(rec) | |
| enriched = self._enrich_with_everything(payload, norm, ts_ms) | |
| self.enriched_signals.append(enriched) | |
| return enriched | |
| def get_recent_enriched(self, ticker: str, timeframe: str, | |
| max_items: int = 300, max_minutes: int = 240) -> List[Dict[str, Any]]: | |
| now_ms = _now_ms() | |
| cutoff = now_ms - max_minutes * 60 * 1000 | |
| out = [ | |
| e for e in reversed(self.enriched_signals) | |
| if e["ticker"] == ticker and str(e["tf"]) == str(timeframe) | |
| and e["timestamp_ms"] >= cutoff | |
| ] | |
| out = list(reversed(out)) | |
| return out[-max_items:] | |
| def _enrich_with_everything(self, payload: Dict[str, Any], norm: Dict[str, Any], ts_ms: int) -> Dict[str, Any]: | |
| technicals = { | |
| "price_metrics": { | |
| "current": payload.get("price"), | |
| "ema_5": payload.get("ema20_5"), | |
| "ema_50": payload.get("ema50_5"), | |
| "sma_20": payload.get("sma20"), | |
| "vwap": payload.get("vwap"), | |
| "bb_upper": payload.get("bb_upper"), | |
| "bb_lower": payload.get("bb_lower"), | |
| "bb_width": payload.get("bb_width"), | |
| "bb_squeeze": payload.get("squeeze"), | |
| "donchian_high": payload.get("don_high"), | |
| "donchian_low": payload.get("don_low"), | |
| "kc_mid": payload.get("kc_mid"), | |
| "hma50": payload.get("hma50"), | |
| "hma_slope": payload.get("hma50_slope"), | |
| }, | |
| "momentum_metrics": { | |
| "rsi": payload.get("rsi"), | |
| "rsi_1h": payload.get("rsi_1h"), | |
| "rsi_4h": payload.get("rsi_4h"), | |
| "stoch_k": payload.get("stoch_k"), | |
| "stoch_d": payload.get("stoch_d"), | |
| "macd_val": payload.get("macd_val"), | |
| "macd_sig": payload.get("macd_sig"), | |
| "adx": payload.get("adx_val"), | |
| "plus_di": payload.get("plus_di"), | |
| "minus_di": payload.get("minus_di"), | |
| "momentum_bullish": payload.get("momentum_bullish"), | |
| "momentum_bearish": payload.get("momentum_bearish"), | |
| }, | |
| "volume_metrics": { | |
| "volume_ratio": payload.get("vol_ratio"), | |
| "volume_zscore": payload.get("vol_zscore"), | |
| "obv_slope": payload.get("obv_slope"), | |
| "cvd_slope": payload.get("cvd_slope"), | |
| "volume_spike": payload.get("volume_spike"), | |
| "volume_oi_ratio": payload.get("volume_oi_ratio"), | |
| }, | |
| "volatility_metrics": { | |
| "atr": payload.get("atr"), | |
| "volatility_ratio": payload.get("volatility_ratio"), | |
| "hurricane_index": payload.get("hurricane_index"), | |
| "market_heat": payload.get("market_heat"), | |
| "impulse": payload.get("impulse"), | |
| }, | |
| "regime_metrics": { | |
| "trend_score": payload.get("trend_score"), | |
| "is_trending": payload.get("is_trending"), | |
| "is_ranging": payload.get("is_ranging"), | |
| "ha_direction": payload.get("ha_dir"), | |
| "don_break_up": payload.get("don_break_up"), | |
| "don_break_dn": payload.get("don_break_dn"), | |
| }, | |
| } | |
| mtf_analysis = { | |
| "mtf_alignment": self._calculate_mtf_alignment(payload), | |
| "higher_tf_bias": self._get_higher_tf_bias(payload), | |
| "timeframe_confluence": self._timeframe_confluence_score(payload), | |
| } | |
| advanced_metrics = { | |
| "composite_momentum": self._composite_momentum_score(technicals), | |
| "volume_pressure": self._volume_pressure_score(technicals), | |
| "volatility_adjusted_strength": self._volatility_adjusted_strength(technicals), | |
| "regime_optimized_signal": self._regime_optimized_signal(technicals, mtf_analysis), | |
| } | |
| data_quality = self._calculate_data_quality_score(payload) | |
| return { | |
| "timestamp_ms": ts_ms, | |
| "ticker": norm.get("ticker"), | |
| "tf": norm.get("tf"), | |
| "emitter": norm.get("emitter"), | |
| "signal": norm.get("signal"), | |
| "strength": norm.get("strength"), | |
| "price": norm.get("price"), | |
| "technical_analysis": technicals, | |
| "mtf_analysis": mtf_analysis, | |
| "advanced_metrics": advanced_metrics, | |
| "data_quality": data_quality, | |
| } | |
| def _calculate_mtf_alignment(self, p: Dict[str, Any]) -> float: | |
| rsi = _to_float(p.get("rsi")) | |
| r1 = _to_float(p.get("rsi_1h")) | |
| r4 = _to_float(p.get("rsi_4h")) | |
| vals = [x for x in (rsi, r1, r4) if x is not None] | |
| if len(vals) < 2: | |
| return 0.5 | |
| same_side = sum( | |
| 1 for v in vals | |
| if (v > 55 and vals[0] > 55) or (v < 45 and vals[0] < 45) | |
| ) | |
| return round(same_side / len(vals), 3) | |
| def _get_higher_tf_bias(self, p: Dict[str, Any]) -> str: | |
| r1 = _to_float(p.get("rsi_1h")) | |
| r4 = _to_float(p.get("rsi_4h")) | |
| vals = [v for v in (r1, r4) if v is not None] | |
| if not vals: | |
| return "NEUTRAL" | |
| avg = float(np.mean(vals)) | |
| if avg > 55: | |
| return "BULLISH" | |
| if avg < 45: | |
| return "BEARISH" | |
| return "NEUTRAL" | |
| def _timeframe_confluence_score(self, p: Dict[str, Any]) -> float: | |
| score = 0.5 | |
| if p.get("is_trending"): | |
| score += 0.1 | |
| if p.get("don_break_up") or p.get("don_break_dn"): | |
| score += 0.1 | |
| return float(np.clip(score, 0.0, 1.0)) | |
| def _composite_momentum_score(self, tech: Dict[str, Any]) -> float: | |
| m = tech["momentum_metrics"] | |
| vals = [] | |
| for k in ("rsi", "rsi_1h", "rsi_4h", "adx"): | |
| v = _to_float(m.get(k)) | |
| if v is not None: | |
| vals.append(v) | |
| if not vals: | |
| return 0.5 | |
| rsi_component = np.mean([v for v in vals if v <= 100]) / 100.0 | |
| adx = _to_float(m.get("adx")) or 20 | |
| adx_component = np.clip(adx / 50.0, 0, 1) | |
| return float(np.clip(0.6 * rsi_component + 0.4 * adx_component, 0, 1)) | |
| def _volume_pressure_score(self, tech: Dict[str, Any]) -> float: | |
| v = tech["volume_metrics"] | |
| ratio = _to_float(v.get("volume_ratio")) or 1.0 | |
| spike = 1.0 if v.get("volume_spike") else 0.0 | |
| z = abs(_to_float(v.get("volume_zscore")) or 0.0) | |
| score = 0.4 * np.clip(ratio / 2.0, 0, 1) + 0.3 * spike + 0.3 * np.clip(z / 3.0, 0, 1) | |
| return float(np.clip(score, 0, 1)) | |
| def _volatility_adjusted_strength(self, tech: Dict[str, Any]) -> float: | |
| vol = tech["volatility_metrics"] | |
| vol_ratio = _to_float(vol.get("volatility_ratio")) or 1.0 | |
| base = np.clip(vol_ratio, 0.5, 2.0) - 0.5 | |
| return float(np.clip(base * 0.7 + 0.3, 0, 1)) | |
| def _regime_optimized_signal(self, tech: Dict[str, Any], mtf: Dict[str, Any]) -> str: | |
| trend_score = _to_float(tech["regime_metrics"].get("trend_score")) or 0.0 | |
| align = mtf.get("mtf_alignment", 0.5) | |
| if trend_score > 0 and align > 0.6: | |
| return "TREND_BULL" | |
| if trend_score < 0 and align > 0.6: | |
| return "TREND_BEAR" | |
| if tech["regime_metrics"].get("is_ranging"): | |
| return "RANGE" | |
| return "MIXED" | |
| def _calculate_data_quality_score(self, payload: Dict[str, Any]) -> float: | |
| keys = ["price", "rsi", "adx_val", "vol_ratio", "atr", "trend_score", "sma20", "ema50_5"] | |
| present = sum(1 for k in keys if payload.get(k) is not None) | |
| return round(present / len(keys), 3) | |
| data_monster = DataMonsterIngestor() | |
| # ========================= | |
| # PATTERN RECOGNITION | |
| # ========================= | |
| class PatternRecognitionEngine: | |
| def __init__(self): | |
| self.patterns_detected = deque(maxlen=1000) | |
| def detect_all_patterns(self, enriched: Dict[str, Any]) -> List[Dict[str, Any]]: | |
| patterns: List[Dict[str, Any]] = [] | |
| patterns.extend(self._detect_price_patterns(enriched)) | |
| patterns.extend(self._detect_volume_patterns(enriched)) | |
| patterns.extend(self._detect_momentum_patterns(enriched)) | |
| patterns.extend(self._detect_volatility_patterns(enriched)) | |
| patterns.extend(self._detect_mtf_patterns(enriched)) | |
| patterns.extend(self._detect_regime_transitions(enriched)) | |
| patterns = sorted(patterns, key=lambda x: x["confidence"], reverse=True) | |
| self.patterns_detected.extend(patterns[:3]) | |
| return patterns | |
| def _detect_price_patterns(self, data: Dict[str, Any]) -> List[Dict[str, Any]]: | |
| out = [] | |
| tech = data["technical_analysis"]["price_metrics"] | |
| reg = data["technical_analysis"]["regime_metrics"] | |
| price = _to_float(tech.get("current")) | |
| ema50 = _to_float(tech.get("ema_50")) | |
| ema5 = _to_float(tech.get("ema_5")) | |
| if reg.get("don_break_up"): | |
| out.append({"type": "BULLISH_BREAKOUT", "confidence": 0.85, | |
| "description": "Price broke above Donchian channel", "category": "price"}) | |
| if reg.get("don_break_dn"): | |
| out.append({"type": "BEARISH_BREAKDOWN", "confidence": 0.85, | |
| "description": "Price broke below Donchian channel", "category": "price"}) | |
| if ema5 and ema50 and ema5 > ema50 and price and price > ema5: | |
| out.append({"type": "BULLISH_EMA_ALIGNMENT", "confidence": 0.8, | |
| "description": "EMA(5) > EMA(50) & price > EMA(5)", "category": "price"}) | |
| return out | |
| def _detect_volume_patterns(self, data: Dict[str, Any]) -> List[Dict[str, Any]]: | |
| out = [] | |
| v = data["technical_analysis"]["volume_metrics"] | |
| ratio = _to_float(v.get("volume_ratio")) or 1.0 | |
| if ratio > 1.5: | |
| out.append({"type": "HIGH_VOLUME", "confidence": min(ratio / 2.5, 1.0), | |
| "description": "Volume above normal", "category": "volume"}) | |
| if v.get("volume_spike"): | |
| out.append({"type": "VOLUME_SPIKE", "confidence": 0.9, | |
| "description": "Volume spike detected", "category": "volume"}) | |
| return out | |
| def _detect_momentum_patterns(self, data: Dict[str, Any]) -> List[Dict[str, Any]]: | |
| out = [] | |
| m = data["technical_analysis"]["momentum_metrics"] | |
| rsi = _to_float(m.get("rsi")) | |
| adx = _to_float(m.get("adx")) | |
| if rsi is not None: | |
| if rsi > 60: | |
| out.append({"type": "RSI_BULLISH", "confidence": np.clip((rsi - 60) / 20, 0, 1), | |
| "description": "RSI bullish zone", "category": "momentum"}) | |
| elif rsi < 40: | |
| out.append({"type": "RSI_BEARISH", "confidence": np.clip((40 - rsi) / 20, 0, 1), | |
| "description": "RSI bearish zone", "category": "momentum"}) | |
| if adx is not None and adx > 20: | |
| out.append({"type": "STRONG_TREND", "confidence": np.clip((adx - 20) / 30, 0, 1), | |
| "description": "ADX strong trend", "category": "momentum"}) | |
| return out | |
| def _detect_volatility_patterns(self, data: Dict[str, Any]) -> List[Dict[str, Any]]: | |
| out = [] | |
| vol = data["technical_analysis"]["volatility_metrics"] | |
| ratio = _to_float(vol.get("volatility_ratio")) or 1.0 | |
| if ratio < 0.7: | |
| out.append({"type": "LOW_VOL_SQUEEZE", "confidence": 0.7, | |
| "description": "Low volatility squeeze", "category": "volatility"}) | |
| if ratio > 1.5: | |
| out.append({"type": "HIGH_VOL_EXPANSION", "confidence": 0.8, | |
| "description": "High volatility expansion", "category": "volatility"}) | |
| return out | |
| def _detect_mtf_patterns(self, data: Dict[str, Any]) -> List[Dict[str, Any]]: | |
| out = [] | |
| mtf = data["mtf_analysis"] | |
| align = mtf.get("mtf_alignment", 0.5) | |
| if align > 0.7: | |
| out.append({"type": "MTF_CONFLUENCE", "confidence": align, | |
| "description": "Multi‑TF alignment", "category": "mtf"}) | |
| return out | |
| def _detect_regime_transitions(self, data: Dict[str, Any]) -> List[Dict[str, Any]]: | |
| out = [] | |
| regime_signal = data["advanced_metrics"].get("regime_optimized_signal") | |
| if regime_signal in ("TREND_BULL", "TREND_BEAR"): | |
| out.append({"type": "TREND_REGIME", "confidence": 0.8, | |
| "description": f"Trend regime: {regime_signal}", "category": "regime"}) | |
| return out | |
| pattern_engine = PatternRecognitionEngine() | |
| # ========================= | |
| # MARKET REGIME / THRESHOLDS | |
| # ========================= | |
| def session_label_utc() -> str: | |
| h = datetime.utcnow().hour | |
| if 0 <= h < 8: | |
| return "ASIA_SESSION" | |
| if 8 <= h < 14: | |
| return "LONDON_SESSION" | |
| if 14 <= h < 22: | |
| return "NEW_YORK_SESSION" | |
| return "OTHER" | |
| class MarketRegimeDetector: | |
| def detect_from_alerts(self, ticker: str, timeframe: str, lookback: int = 120) -> Dict[str, Any]: | |
| pairs = [r for r in ALERTS if r.get("ticker") == ticker and str(r.get("tf")) == str(timeframe) | |
| and r.get("price") is not None] | |
| if not pairs: | |
| return {"regime": "NORMAL", "volatility": 0.0, "trend_strength": 0.0} | |
| prices = [float(r["price"]) for r in pairs][-lookback:] | |
| if len(prices) < 10: | |
| return {"regime": "NORMAL", "volatility": 0.0, "trend_strength": 0.0} | |
| arr = np.array(prices, dtype=float) | |
| vol = float(np.std(arr[-20:]) / max(np.mean(arr[-20:]), 1e-9)) | |
| t = np.arange(len(arr)) | |
| corr = float(np.corrcoef(t, arr)[0, 1]) if len(arr) > 1 else 0.0 | |
| ts = abs(corr) | |
| if vol > 0.02 and ts < 0.3: | |
| regime = "HIGH_VOLATILITY_RANGING" | |
| elif vol > 0.02 and ts > 0.7: | |
| regime = "HIGH_VOLATILITY_TRENDING" | |
| elif vol < 0.01 and ts > 0.7: | |
| regime = "LOW_VOLATILITY_TRENDING" | |
| else: | |
| regime = "NORMAL" | |
| return {"regime": regime, "volatility": round(vol, 4), "trend_strength": round(ts, 3)} | |
| def optimize_alerts_threshold(signal_strength: int, market_regime: str, time_of_day: str) -> Tuple[bool, int]: | |
| base_threshold = 70 | |
| adjustments = { | |
| "HIGH_VOLATILITY_RANGING": +15, | |
| "LOW_VOLATILITY_TRENDING": -10, | |
| "ASIA_SESSION": +5, | |
| "LONDON_SESSION": -5 | |
| } | |
| thr = base_threshold + adjustments.get(market_regime, 0) + adjustments.get(time_of_day, 0) | |
| return signal_strength >= thr, thr | |
| regime_detector = MarketRegimeDetector() | |
| # ========================= | |
| # FUSION ENGINE (base + enhanced) | |
| # ========================= | |
| class SignalFusionEngine: | |
| def calculate_fusion_signal(self, ticker: str, timeframe: str) -> Dict[str, Any]: | |
| cfg = RUNTIME_CONFIG["fusion"] | |
| window_ms = cfg["time_window_minutes"] * 60 * 1000 | |
| current_time = _now_ms() | |
| recent = [a for a in ALERTS | |
| if a.get("ticker") == ticker and str(a.get("tf")) == str(timeframe) | |
| and current_time - a["received_ms"] <= window_ms] | |
| if not recent: | |
| return self._create_no_signal_response(ticker, timeframe) | |
| source_signals: Dict[str, Dict[str, Any]] = {} | |
| for s in recent: | |
| source_signals[s.get("emitter", "unknown")] = s | |
| signal_values = [] | |
| total_weight = 0.0 | |
| details = [] | |
| for src, sig in source_signals.items(): | |
| w = float(source_perf.get_source_weight(src)) | |
| sig_txt = (sig.get("signal") or "FLAT").upper() | |
| numeric = 1.0 if sig_txt.startswith("LONG") else (-1.0 if sig_txt.startswith("SHORT") else 0.0) | |
| strength = float(sig.get("strength", 50)) / 100.0 | |
| adjusted_w = w * strength | |
| if adjusted_w >= cfg["weight_threshold"]: | |
| contrib = numeric * adjusted_w | |
| signal_values.append(contrib) | |
| total_weight += adjusted_w | |
| details.append({ | |
| "source": src, | |
| "signal": sig_txt, | |
| "strength": int(sig.get("strength", 0)), | |
| "weight": round(adjusted_w, 3), | |
| "contribution": round(contrib, 4) | |
| }) | |
| if total_weight == 0: | |
| return self._create_weak_signal_response(ticker, timeframe, details) | |
| fusion_value = sum(signal_values) / total_weight | |
| if fusion_value >= 0.25: | |
| final_signal = "STRONG LONG" | |
| confidence = min(fusion_value * 2, 1.0) | |
| elif fusion_value >= 0.10: | |
| final_signal = "MODERATE LONG" | |
| confidence = fusion_value * 3 | |
| elif fusion_value <= -0.25: | |
| final_signal = "STRONG SHORT" | |
| confidence = min(abs(fusion_value) * 2, 1.0) | |
| elif fusion_value <= -0.10: | |
| final_signal = "MODERATE SHORT" | |
| confidence = abs(fusion_value) * 3 | |
| else: | |
| final_signal = "NEUTRAL" | |
| confidence = 0.5 - abs(fusion_value) | |
| if confidence < cfg["min_confidence"]: | |
| final_signal = "WEAK SIGNAL" | |
| consensus = self._consensus(signal_values) | |
| return { | |
| "ticker": ticker, | |
| "timeframe": str(timeframe), | |
| "fusion_signal": final_signal, | |
| "fusion_confidence": round(confidence, 3), | |
| "fusion_score": round(fusion_value, 3), | |
| "consensus": round(consensus, 3), | |
| "sources_used": len(details), | |
| "source_details": details, | |
| "timestamp": datetime.now(tz=ATHENS_TZ).isoformat(), | |
| "recommendation": self._recommendation(final_signal, confidence, consensus, cfg["consensus_threshold"]) | |
| } | |
| def _consensus(values: List[float]) -> float: | |
| if not values: | |
| return 0.0 | |
| long_c = sum(1 for v in values if v > 0) | |
| short_c = sum(1 for v in values if v < 0) | |
| neutral_c = sum(1 for v in values if v == 0) | |
| total = len(values) | |
| return max(long_c, short_c, neutral_c) / total if total else 0.0 | |
| def _recommendation(signal: str, conf: float, cons: float, cons_thr: float) -> str: | |
| base = { | |
| "STRONG LONG": ["🎯 ΙΣΧΥΡΟ LONG - Άμεση είσοδος", "📈 Υψηλή πιθανότητα ανόδου", "💪 Πολλαπλές επιβεβαιώσεις"], | |
| "MODERATE LONG": ["✅ MODERATE LONG - Είσοδος με προσοχή", "📊 Αναμενόμενη άνοδος"], | |
| "STRONG SHORT": ["🎯 ΙΣΧΥΡΟ SHORT - Άμεση είσοδος", "📉 Υψηλή πιθανότητα πτώσης", "💪 Πολλαπλές επιβεβαιώσεις"], | |
| "MODERATE SHORT": ["✅ MODERATE SHORT - Είσοδος με προσοχή", "📊 Αναμενόμενη πτώση"], | |
| "NEUTRAL": ["⚪ NEUTRAL - Αποφυγή εισόδων", "🔍 Περιμένουμε καλύτερο σήμα"], | |
| "WEAK SIGNAL": ["⚠️ Ασθενές σήμα - Μη αξιοποιήσιμο", "🔍 Περιμένουμε επιβεβαιώσεις"] | |
| }.get(signal, ["📊 Αναμονή για σήμα..."]) | |
| if conf >= 0.8: | |
| base.append("🔥 Πολύ υψηλή εμπιστοσύνη") | |
| elif conf >= 0.6: | |
| base.append("👍 Καλή εμπιστοσύνη") | |
| else: | |
| base.append("⚠️ Χαμηλή εμπιστοσύνη") | |
| if cons >= 0.8: | |
| base.append("🤝 Πολύ υψηλή ομοφωνία") | |
| elif cons >= cons_thr: | |
| base.append("👥 Καλή ομοφωνία") | |
| return " | ".join(base) | |
| def _create_no_signal_response(ticker: str, timeframe: str) -> Dict[str, Any]: | |
| return { | |
| "ticker": ticker, | |
| "timeframe": timeframe, | |
| "fusion_signal": "NO SIGNAL", | |
| "fusion_confidence": 0.0, | |
| "fusion_score": 0.0, | |
| "consensus": 0.0, | |
| "sources_used": 0, | |
| "source_details": [], | |
| "timestamp": datetime.now(tz=ATHENS_TZ).isoformat(), | |
| "recommendation": "🔍 Δεν υπάρχουν πρόσφατα σήματα" | |
| } | |
| def _create_weak_signal_response(ticker: str, timeframe: str, details: List[Dict[str, Any]]) -> Dict[str, Any]: | |
| return { | |
| "ticker": ticker, | |
| "timeframe": timeframe, | |
| "fusion_signal": "WEAK SIGNAL", | |
| "fusion_confidence": 0.0, | |
| "fusion_score": 0.0, | |
| "consensus": 0.0, | |
| "sources_used": len(details), | |
| "source_details": details, | |
| "timestamp": datetime.now(tz=ATHENS_TZ).isoformat(), | |
| "recommendation": "⚠️ Πολύ ασθενή σήματα - Αποφυγή εισόδων" | |
| } | |
| class EnhancedSignalFusionEngine(SignalFusionEngine): | |
| def calculate_enhanced_fusion(self, ticker: str, timeframe: str) -> Dict[str, Any]: | |
| base = super().calculate_fusion_signal(ticker, timeframe) | |
| momentum_score = self._momentum_from_alerts(ticker, timeframe) | |
| volume_score = self._volume_from_alerts(ticker, timeframe) | |
| enhanced_conf = base["fusion_confidence"] * 0.7 + momentum_score * 0.2 + volume_score * 0.1 | |
| strengths = [d["strength"] for d in base.get("source_details", [])] or [0] | |
| avg_strength = (sum(strengths) / len(strengths)) / 100.0 | |
| signals = [d["signal"] for d in base.get("source_details", [])] | |
| consensus = (max(signals.count(s) for s in set(signals)) / len(signals)) if signals else 0.0 | |
| quality = consensus * 0.6 + avg_strength * 0.4 | |
| regime_info = regime_detector.detect_from_alerts(ticker, str(timeframe)) | |
| session = session_label_utc() | |
| trigger, dyn_thr = optimize_alerts_threshold(int(avg_strength * 100), regime_info["regime"], session) | |
| base.update({ | |
| "enhanced_confidence": round(enhanced_conf, 3), | |
| "momentum_score": round(momentum_score, 3), | |
| "volume_score": round(volume_score, 3), | |
| "quality_score": round(quality, 3), | |
| "composite_score": round(enhanced_conf * quality, 3), | |
| "market_regime": regime_info["regime"], | |
| "regime_volatility": regime_info["volatility"], | |
| "regime_trend_strength": regime_info["trend_strength"], | |
| "session": session, | |
| "dynamic_threshold": dyn_thr, | |
| "trigger": trigger | |
| }) | |
| return base | |
| def _momentum_from_alerts(ticker: str, timeframe: str, lookback: int = 20) -> float: | |
| vals = [r.get("price") for r in ALERTS | |
| if r.get("ticker") == ticker and str(r.get("tf")) == str(timeframe) and r.get("price") is not None] | |
| if len(vals) < 5: | |
| return 0.5 | |
| arr = np.array(vals[-lookback:], dtype=float) | |
| t = np.arange(len(arr)) | |
| slope = np.polyfit(t, arr, 1)[0] | |
| score = 0.5 + np.tanh(slope / (np.std(arr) + 1e-9)) * 0.5 | |
| return float(np.clip(score, 0.0, 1.0)) | |
| def _volume_from_alerts(ticker: str, timeframe: str, lookback: int = 30) -> float: | |
| vols = [] | |
| for r in reversed(ALERTS): | |
| if r.get("ticker") == ticker and str(r.get("tf")) == str(timeframe): | |
| raw = r.get("raw") or {} | |
| vr = _to_float(raw.get("vol_ratio")) | |
| if vr is not None: | |
| vols.append(vr) | |
| if len(vols) >= lookback: | |
| break | |
| if not vols: | |
| return 0.5 | |
| v = float(np.mean(vols)) | |
| if v >= 2.0: | |
| return 1.0 | |
| if v >= 1.5: | |
| return 0.7 | |
| return 0.4 | |
| enhanced_engine = EnhancedSignalFusionEngine() | |
| # ========================= | |
| # PRO+ LAYER (validator, persistence, risk-adjusted) | |
| # ========================= | |
| class SmartSignalValidator: | |
| def validate_signal_quality(self, fusion_result: Dict[str, Any], ticker: str, timeframe: str) -> Dict[str, Any]: | |
| volume_ok = self._check_volume_support(ticker) | |
| price_action_ok = self._check_price_action_alignment(fusion_result, ticker, timeframe) | |
| time_effectiveness = self._session_effectiveness_score() | |
| persistence_score = self._signal_persistence(ticker, timeframe, fusion_result["fusion_signal"]) | |
| quality_score = ( | |
| volume_ok * 0.3 + | |
| price_action_ok * 0.3 + | |
| time_effectiveness * 0.2 + | |
| persistence_score * 0.2 | |
| ) | |
| if quality_score < 0.6: | |
| fusion_result["fusion_signal"] = "WEAK SIGNAL" | |
| fusion_result["recommendation"] += " | ⚠️ Χαμηλή Ποιότητα Σήματος" | |
| fusion_result["quality_metrics"] = { | |
| "volume_confirmation": volume_ok, | |
| "price_action_alignment": price_action_ok, | |
| "session_effectiveness": time_effectiveness, | |
| "signal_persistence": persistence_score, | |
| "overall_quality": round(quality_score, 3) | |
| } | |
| return fusion_result | |
| def _check_volume_support(self, ticker: str) -> float: | |
| recent_alerts = [r for r in ALERTS if r.get("ticker") == ticker][-20:] | |
| if not recent_alerts: | |
| return 0.5 | |
| volumes = [] | |
| for alert in recent_alerts: | |
| raw = alert.get("raw", {}) | |
| vol_ratio = raw.get("vol_ratio") | |
| if vol_ratio is not None: | |
| volumes.append(float(vol_ratio)) | |
| if not volumes: | |
| return 0.5 | |
| avg_volume = sum(volumes) / len(volumes) | |
| return 0.8 if avg_volume > 1.2 else 0.5 if avg_volume > 0.8 else 0.3 | |
| def _check_price_action_alignment(self, fusion_result: Dict, ticker: str, timeframe: str) -> float: | |
| recent_alerts = [r for r in ALERTS | |
| if r.get("ticker") == ticker and str(r.get("tf")) == str(timeframe)][-10:] | |
| if not recent_alerts: | |
| return 0.5 | |
| prices = [r.get("price") for r in recent_alerts if r.get("price") is not None] | |
| if len(prices) < 5: | |
| return 0.5 | |
| price_trend = np.polyfit(range(len(prices)), prices, 1)[0] | |
| sig = (fusion_result["fusion_signal"] or "").upper() | |
| if "LONG" in sig and price_trend > 0: | |
| return 0.8 | |
| if "SHORT" in sig and price_trend < 0: | |
| return 0.8 | |
| return 0.3 | |
| def _session_effectiveness_score(self) -> float: | |
| sess = session_label_utc() | |
| if sess == "LONDON_SESSION": | |
| return 0.8 | |
| if sess == "NEW_YORK_SESSION": | |
| return 0.7 | |
| return 0.5 | |
| _persistence_store = defaultdict(lambda: deque(maxlen=10)) | |
| def _signal_persistence(self, ticker: str, timeframe: str, current_signal: str) -> float: | |
| key = f"{ticker}|{timeframe}" | |
| dq = self._persistence_store[key] | |
| dq.append(current_signal) | |
| trends = list(dq) | |
| if len(trends) < 3: | |
| return 0.5 | |
| unique_signals = len(set(trends)) | |
| if unique_signals == 1: | |
| return 0.9 | |
| if unique_signals == 2: | |
| return 0.7 | |
| return 0.4 | |
| def calculate_dynamic_confidence(fusion_result: Dict, market_regime: str) -> float: | |
| base_confidence = fusion_result["fusion_confidence"] | |
| sources_used = fusion_result["sources_used"] | |
| consensus = fusion_result["consensus"] | |
| regime_weights = { | |
| "HIGH_VOLATILITY_RANGING": 0.7, | |
| "LOW_VOLATILITY_TRENDING": 1.1, | |
| "HIGH_VOLATILITY_TRENDING": 0.9, | |
| "NORMAL": 1.0 | |
| } | |
| regime_multiplier = regime_weights.get(market_regime, 1.0) | |
| diversity_bonus = min(sources_used / 5, 0.3) | |
| consensus_strength = consensus * 0.2 | |
| adjusted_confidence = base_confidence * regime_multiplier + diversity_bonus + consensus_strength | |
| return float(min(adjusted_confidence, 1.0)) | |
| class RealTimePerformanceMonitor: | |
| def __init__(self): | |
| self.performance_log = deque(maxlen=1000) | |
| def log_trade_signal(self, signal_data: Dict, actual_outcome: Optional[float] = None): | |
| self.performance_log.append({ | |
| "timestamp": datetime.now(tz=ATHENS_TZ).isoformat(), | |
| "signal": signal_data.get("fusion_signal"), | |
| "confidence": signal_data.get("fusion_confidence"), | |
| "sources": signal_data.get("sources_used"), | |
| "ticker": signal_data.get("ticker"), | |
| "timeframe": signal_data.get("timeframe"), | |
| "regime": signal_data.get("market_regime", "UNKNOWN"), | |
| "actual_outcome": actual_outcome, | |
| "quality_metrics": signal_data.get("quality_metrics", {}) | |
| }) | |
| def calculate_real_time_metrics(self) -> Dict[str, Any]: | |
| if not self.performance_log: | |
| return {} | |
| recent = list(self.performance_log)[-100:] | |
| total = len(recent) | |
| winning = sum(1 for s in recent if (s.get("actual_outcome") or 0) > 0) | |
| win_rate = winning / total if total else 0 | |
| high_conf = [s for s in recent if (s.get("confidence") or 0) > 0.7] | |
| hc_win = sum(1 for s in high_conf if (s.get("actual_outcome") or 0) > 0) | |
| hc_rate = hc_win / len(high_conf) if high_conf else 0 | |
| avg_conf = sum((s.get("confidence") or 0) for s in recent) / total if total else 0 | |
| return { | |
| "recent_win_rate": round(win_rate * 100, 2), | |
| "high_confidence_accuracy": round(hc_rate * 100, 2), | |
| "signals_tracked": total, | |
| "avg_confidence": round(avg_conf, 3) | |
| } | |
| def generate_risk_adjusted_recommendation(fusion_result: Dict, account_size: float = 10_000) -> Dict[str, Any]: | |
| signal = fusion_result["fusion_signal"] | |
| confidence = fusion_result["fusion_confidence"] | |
| quality_metrics = fusion_result.get("quality_metrics", {}) | |
| base_recommendation = fusion_result["recommendation"] | |
| if "STRONG" in signal and confidence > 0.8 and quality_metrics.get("overall_quality", 0) > 0.7: | |
| risk_level = "LOW" | |
| position_size = "3-5%" | |
| elif "MODERATE" in signal and confidence > 0.6: | |
| risk_level = "MEDIUM" | |
| position_size = "1-3%" | |
| else: | |
| risk_level = "HIGH" | |
| position_size = "0.5-1%" | |
| regime = fusion_result.get("market_regime", "NORMAL") | |
| if regime == "HIGH_VOLATILITY_RANGING": | |
| position_size = "0.5-1%" | |
| risk_level = "HIGH" | |
| return { | |
| "signal": signal, | |
| "confidence": confidence, | |
| "risk_level": risk_level, | |
| "recommended_position_size": position_size, | |
| "entry_timing": "IMMEDIATE" if "STRONG" in signal else "WAIT_FOR_CONFIRMATION", | |
| "stop_loss_advice": "TIGHT" if regime == "HIGH_VOLATILITY_RANGING" else "STANDARD", | |
| "detailed_reasoning": base_recommendation, | |
| "quality_score": quality_metrics.get("overall_quality", 0.5) | |
| } | |
| class UltraXProPlusEngine(EnhancedSignalFusionEngine): | |
| def __init__(self): | |
| super().__init__() | |
| self.validator = SmartSignalValidator() | |
| self.performance_monitor = RealTimePerformanceMonitor() | |
| def calculate_pro_plus_fusion(self, ticker: str, timeframe: str) -> Dict[str, Any]: | |
| base = super().calculate_enhanced_fusion(ticker, timeframe) | |
| validated = self.validator.validate_signal_quality(base, ticker, timeframe) | |
| regime = validated.get("market_regime", "NORMAL") | |
| dynamic_confidence = calculate_dynamic_confidence(validated, regime) | |
| validated["dynamic_confidence"] = round(dynamic_confidence, 3) | |
| risk_adjusted = generate_risk_adjusted_recommendation(validated) | |
| validated["risk_adjusted"] = risk_adjusted | |
| self.performance_monitor.log_trade_signal(validated) | |
| return validated | |
| pro_plus_engine = UltraXProPlusEngine() | |
| # ========================= | |
| # MEGA FUSION (DataMonster + patterns) | |
| # ========================= | |
| class MegaFusionEngineV2: | |
| def __init__(self, ingestor: DataMonsterIngestor, pattern_engine: PatternRecognitionEngine): | |
| self.ingestor = ingestor | |
| self.pattern_engine = pattern_engine | |
| def process_mega_fusion(self, ticker: str, timeframe: str) -> Dict[str, Any]: | |
| base = pro_plus_engine.calculate_pro_plus_fusion(ticker, timeframe) | |
| recent = self.ingestor.get_recent_enriched(ticker, timeframe, max_items=200, max_minutes=240) | |
| if not recent: | |
| base["mega_info"] = {"reason": "no_enriched_data"} | |
| return base | |
| all_patterns: List[Dict[str, Any]] = [] | |
| for e in recent: | |
| all_patterns.extend(self.pattern_engine.detect_all_patterns(e)) | |
| data_quality_vals = [float(e.get("data_quality", 0.5)) for e in recent] | |
| avg_data_quality = float(np.mean(data_quality_vals)) if data_quality_vals else 0.5 | |
| long_cnt = sum(1 for e in recent if str(e.get("signal","")).upper().startswith("LONG")) | |
| short_cnt = sum(1 for e in recent if str(e.get("signal","")).upper().startswith("SHORT")) | |
| flat_cnt = len(recent) - long_cnt - short_cnt | |
| total = max(len(recent), 1) | |
| direction_consistency = max(long_cnt, short_cnt, flat_cnt) / total | |
| bullish_patterns = sum(1 for p in all_patterns if "BULL" in p["type"]) | |
| bearish_patterns = sum(1 for p in all_patterns if "BEAR" in p["type"]) | |
| pattern_bias = (bullish_patterns - bearish_patterns) / max(len(all_patterns), 1) | |
| validation_score = float(np.clip( | |
| 0.4 * direction_consistency + | |
| 0.3 * avg_data_quality + | |
| 0.3 * abs(pattern_bias), | |
| 0, 1 | |
| )) | |
| dyn_conf = float(base.get("dynamic_confidence", base.get("fusion_confidence", 0.0))) | |
| mega_conf = float(np.clip( | |
| 0.5 * dyn_conf + | |
| 0.3 * validation_score + | |
| 0.2 * avg_data_quality, | |
| 0, 1 | |
| )) | |
| base_signal = base.get("fusion_signal", "NEUTRAL") | |
| if "STRONG" in base_signal and mega_conf > 0.8: | |
| mega_label = "MEGA " + base_signal | |
| elif "MODERATE" in base_signal and mega_conf > 0.7: | |
| mega_label = "MEGA " + base_signal | |
| else: | |
| mega_label = base_signal | |
| mega_recommendation = base.get("risk_adjusted", {}).get("detailed_reasoning", base.get("recommendation","")) | |
| mega_recommendation += f" | 🧠 Mega validation: {validation_score:.2f}, dataQ: {avg_data_quality:.2f}" | |
| base.update({ | |
| "fusion_signal": mega_label, | |
| "fusion_confidence": round(mega_conf, 3), | |
| "mega_validation_score": round(validation_score, 3), | |
| "mega_data_quality": round(avg_data_quality, 3), | |
| "mega_patterns_count": len(all_patterns), | |
| "mega_direction_consistency": round(direction_consistency, 3), | |
| "mega_pattern_bias": round(pattern_bias, 3), | |
| "recommendation": mega_recommendation, | |
| }) | |
| return base | |
| mega_engine = MegaFusionEngineV2(data_monster, pattern_engine) | |
| # ========================= | |
| # STORAGE / ALERTS / WS | |
| # ========================= | |
| ALERTS: deque = deque(maxlen=MAX_ALERTS) | |
| FUSION_HISTORY: deque = deque(maxlen=500) | |
| class WSManager: | |
| def __init__(self): | |
| self.conns: List[WebSocket] = [] | |
| async def connect(self, ws: WebSocket): | |
| await ws.accept() | |
| self.conns.append(ws) | |
| def disconnect(self, ws: WebSocket): | |
| if ws in self.conns: | |
| self.conns.remove(ws) | |
| async def broadcast(self, payload: Dict[str, Any]): | |
| dead = [] | |
| for c in list(self.conns): | |
| try: | |
| await c.send_json(payload) | |
| except Exception: | |
| dead.append(c) | |
| for c in dead: | |
| self.disconnect(c) | |
| ws_manager = WSManager() | |
| def log_fusion_result(result: Dict[str, Any], mode: str): | |
| try: | |
| row = { | |
| "ts_ms": _now_ms(), | |
| "when": athens_str_from_ms(_now_ms()), | |
| "mode": mode, | |
| "ticker": result.get("ticker"), | |
| "tf": result.get("timeframe"), | |
| "signal": result.get("fusion_signal"), | |
| "conf%": int(round(result.get("fusion_confidence", 0) * 100)), | |
| "dyn_conf%": int(round(result.get("dynamic_confidence", result.get("enhanced_confidence", 0)) * 100)), | |
| "sources": result.get("sources_used", 0), | |
| "regime": result.get("market_regime", ""), | |
| "session": result.get("session", "") | |
| } | |
| FUSION_HISTORY.append(row) | |
| except Exception: | |
| pass | |
| def add_alert(p: Dict[str, Any]) -> Dict[str, Any]: | |
| norm = normalize_payload(p) | |
| try: | |
| data_monster.ingest_from_normalized(norm) | |
| except Exception: | |
| pass | |
| row = { | |
| "received_ms": _now_ms(), | |
| "schema": norm["schema"], | |
| "emitter": norm["emitter"], | |
| "ticker": norm["ticker"], | |
| "tf": norm["tf"], | |
| "signal": norm["signal"], | |
| "strength": norm["strength"], | |
| "price": norm["price"], | |
| "SL": norm["sl"], | |
| "TP": norm["tp"], | |
| "RR": norm["rr_ratio"], | |
| "pos%": norm["position_size"], | |
| "ADX": norm["adx"], | |
| "RSI": norm["rsi"], | |
| "bar_time": norm["bar_time"] or ms_to_iso_utc(norm["bar_time_ms"]) or "", | |
| "raw": norm["raw"], | |
| } | |
| ALERTS.append(row) | |
| persist_append(row) | |
| try: | |
| import asyncio | |
| loop = asyncio.get_running_loop() | |
| pub = {k: v for k, v in row.items() if k != "raw"} | |
| loop.create_task(ws_manager.broadcast({"type": "new_alert", "data": pub})) | |
| except RuntimeError: | |
| pass | |
| return row | |
| # ========================= | |
| # STATS / TABLES / HISTORY | |
| # ========================= | |
| def latest_table(ticker_filter: Optional[str] = None, limit: int = 1000) -> pd.DataFrame: | |
| rows: List[Dict[str, Any]] = list(ALERTS)[-limit:] | |
| if ticker_filter: | |
| t = ticker_filter.strip().upper() | |
| rows = [r for r in rows if (r.get("ticker", "") or "").upper().startswith(t)] | |
| latest: Dict[str, Dict[str, Any]] = {} | |
| for r in rows: | |
| key = f'{r.get("emitter")}|{r.get("ticker")}|{r.get("tf")}' | |
| if key not in latest or r["received_ms"] > latest[key]["received_ms"]: | |
| latest[key] = r | |
| recs = sorted(latest.values(), key=lambda x: x["received_ms"], reverse=True) | |
| def sig_mark(s: str) -> str: | |
| s = (s or "").upper() | |
| if s.startswith("LONG"): | |
| return "🟢 LONG" | |
| if s.startswith("SHORT"): | |
| return "🔴 SHORT" | |
| return "⚪ FLAT" | |
| return pd.DataFrame([{ | |
| "when": athens_str_from_ms(r["received_ms"]), | |
| "src": r.get("emitter"), | |
| "ticker": r.get("ticker"), | |
| "tf": r.get("tf"), | |
| "signal": sig_mark(r.get("signal")), | |
| "strength": r.get("strength"), | |
| "price": r.get("price"), | |
| "SL": r.get("SL"), | |
| "TP": r.get("TP"), | |
| "RR": r.get("RR"), | |
| "pos%": r.get("pos%"), | |
| "ADX": r.get("ADX"), | |
| "RSI": r.get("RSI"), | |
| "bar_time": r.get("bar_time"), | |
| } for r in recs]) | |
| def compute_stats() -> Dict[str, Any]: | |
| rows = list(ALERTS) | |
| emitters: Dict[str, Dict[str, Any]] = {} | |
| all_tickers = set() | |
| for r in rows: | |
| e = r.get("emitter") or "unknown" | |
| d = emitters.setdefault(e, {"count": 0, "tickers": set(), "tfs": set(), | |
| "last_received_ms": 0, "last_bar_time": ""}) | |
| d["count"] += 1 | |
| if r.get("ticker"): | |
| d["tickers"].add(r["ticker"]) | |
| if r.get("tf"): | |
| d["tfs"].add(r["tf"]) | |
| if r["received_ms"] > d["last_received_ms"]: | |
| d["last_received_ms"] = r["received_ms"] | |
| d["last_bar_time"] = r.get("bar_time") or d["last_bar_time"] | |
| if r.get("ticker"): | |
| all_tickers.add(r["ticker"]) | |
| per = [] | |
| for e, d in emitters.items(): | |
| per.append({ | |
| "src": e, | |
| "signals": d["count"], | |
| "tickers": len([x for x in d["tickers"] if x]), | |
| "tfs": len([x for x in d["tfs"] if x]), | |
| "last_seen": athens_str_from_ms(d["last_received_ms"]) if d["last_received_ms"] else "", | |
| "last_bar_time": d["last_bar_time"], | |
| }) | |
| per.sort(key=lambda x: (x["signals"], x["last_seen"]), reverse=True) | |
| def _active_count(minutes: int) -> int: | |
| cut = _now_ms() - minutes * 60 * 1000 | |
| return sum(1 for r in rows if r["received_ms"] >= cut) | |
| return { | |
| "total_events": len(rows), | |
| "unique_emitters": len(emitters), | |
| "unique_tickers": len([x for x in all_tickers if x]), | |
| "active_30": _active_count(30), | |
| "active_5": _active_count(5), | |
| "per_emitter": per, | |
| } | |
| def active_sources(minutes: int = 60, ticker: Optional[str] = None, tf: Optional[str] = None) -> pd.DataFrame: | |
| cut = _now_ms() - minutes * 60 * 1000 | |
| rows = [r for r in ALERTS if r["received_ms"] >= cut] | |
| if ticker: | |
| t = ticker.upper() | |
| rows = [r for r in rows if (r.get("ticker", "") or "").upper() == t] | |
| if tf: | |
| rows = [r for r in rows if str(r.get("tf")) == str(tf)] | |
| agg: Dict[str, Dict[str, Any]] = {} | |
| for r in rows: | |
| key = f'{r.get("ticker")}|{r.get("tf")}' | |
| a = agg.setdefault(key, {"ticker": r.get("ticker"), "tf": r.get("tf"), | |
| "sources": set(), "last_seen_ms": 0}) | |
| a["sources"].add(r.get("emitter")) | |
| if r["received_ms"] > a["last_seen_ms"]: | |
| a["last_seen_ms"] = r["received_ms"] | |
| out = [] | |
| for a in agg.values(): | |
| out.append({ | |
| "ticker": a["ticker"], | |
| "tf": a["tf"], | |
| "sources_count": len(a["sources"]), | |
| "sources": ", ".join(sorted(a["sources"])), | |
| "last_seen": athens_str_from_ms(a["last_seen_ms"]) if a["last_seen_ms"] else "" | |
| }) | |
| out.sort(key=lambda x: (x["sources_count"], x["last_seen"]), reverse=True) | |
| return pd.DataFrame(out) | |
| def fusion_history_df(limit: int = 50) -> pd.DataFrame: | |
| rows = list(FUSION_HISTORY)[-limit:] | |
| rows = sorted(rows, key=lambda x: x["ts_ms"], reverse=True) | |
| return pd.DataFrame([{ | |
| "when": r["when"], | |
| "mode": r["mode"], | |
| "ticker": r["ticker"], | |
| "tf": r["tf"], | |
| "signal": r["signal"], | |
| "conf%": r["conf%"], | |
| "dyn_conf%": r["dyn_conf%"], | |
| "sources": r["sources"], | |
| "regime": r["regime"], | |
| "session": r["session"], | |
| } for r in rows]) | |
| # ========================= | |
| # FASTAPI | |
| # ========================= | |
| api = FastAPI(title="ULTRA‑X Multi‑Source API", docs_url="/api/docs", openapi_url="/api/openapi.json") | |
| api.add_middleware(CORSMiddleware, allow_origins=["*"], allow_credentials=False, allow_methods=["*"], allow_headers=["*"]) | |
| def healthz(): | |
| return {"ok": True, "alerts_cached": len(ALERTS), "persist": bool(PERSIST_PATH)} | |
| def index(): | |
| return { | |
| "ok": True, | |
| "ui": "/", | |
| "endpoints": [ | |
| "POST /api/webhook?token=YOUR_TOKEN", | |
| "GET /api/latest", | |
| "GET /api/stats", | |
| "GET /api/fusion?ticker=BTCUSD&timeframe=5", | |
| "GET /api/fusion/enhanced?ticker=BTCUSD&timeframe=5", | |
| "GET /api/fusion/proplus?ticker=BTCUSD&timeframe=5", | |
| "GET /api/fusion/mega?ticker=BTCUSD&timeframe=5", | |
| "GET /api/sources/active?minutes=5&ticker=BTCUSD&tf=5", | |
| "GET /api/config", | |
| "POST /api/config", | |
| "POST /api/perf/update", | |
| "GET /api/perf/scores", | |
| "GET /api/market/context?ticker=BTCUSD&timeframe=5", | |
| "GET /api/docs", | |
| "WS /ws" | |
| ], | |
| } | |
| def _check_body_secret(payload: Dict[str, Any]): | |
| if not BODY_SECRETS: | |
| return | |
| if "secret" in payload and payload["secret"] not in BODY_SECRETS: | |
| raise HTTPException(status_code=401, detail="Invalid body secret") | |
| async def webhook(request: Request, token: Optional[str] = Query(default=None)): | |
| if WEBHOOK_TOKEN and token != WEBHOOK_TOKEN: | |
| raise HTTPException(status_code=401, detail="Invalid token") | |
| try: | |
| payload = await request.json() | |
| except Exception: | |
| raise HTTPException(status_code=400, detail="Invalid JSON body") | |
| if isinstance(payload, dict): | |
| _check_body_secret(payload) | |
| row = add_alert(payload) | |
| return {"ok": True, "result": {k: v for k, v in row.items() if k != "raw"}} | |
| elif isinstance(payload, list): | |
| count = 0 | |
| for p in payload: | |
| if isinstance(p, dict): | |
| _check_body_secret(p) | |
| add_alert(p) | |
| count += 1 | |
| return {"ok": True, "count": count} | |
| else: | |
| raise HTTPException(status_code=400, detail="Unsupported payload type") | |
| def latest(limit: int = 1000, ticker: Optional[str] = None): | |
| df = latest_table(ticker_filter=ticker, limit=limit) | |
| return {"ok": True, "items": df.to_dict(orient="records")} | |
| def export_csv(limit: int = 5000, ticker: Optional[str] = None): | |
| df = latest_table(ticker_filter=ticker, limit=limit) | |
| csv = df.to_csv(index=False) | |
| return StreamingResponse(iter([csv]), media_type="text/csv", | |
| headers={"Content-Disposition": "attachment; filename=signals.csv"}) | |
| def stats(): | |
| s = compute_stats() | |
| return {"ok": True, **s} | |
| def get_fusion_signal(ticker: str = Query("BTCUSD"), timeframe: str = Query("5")): | |
| base = enhanced_engine.calculate_fusion_signal(ticker, timeframe) | |
| return {"ok": True, "fusion": base} | |
| def get_enhanced_fusion(ticker: str = Query("BTCUSD"), timeframe: str = Query("5")): | |
| result = enhanced_engine.calculate_enhanced_fusion(ticker, timeframe) | |
| return {"ok": True, "fusion": result} | |
| def get_proplus_fusion(ticker: str = Query("BTCUSD"), timeframe: str = Query("5")): | |
| result = pro_plus_engine.calculate_pro_plus_fusion(ticker, timeframe) | |
| return {"ok": True, "fusion": result} | |
| def get_mega_fusion(ticker: str = Query("BTCUSD"), timeframe: str = Query("5")): | |
| result = mega_engine.process_mega_fusion(ticker, timeframe) | |
| return {"ok": True, "fusion": result} | |
| def get_active_sources(minutes: int = 60, ticker: Optional[str] = None, tf: Optional[str] = None): | |
| df = active_sources(minutes, ticker, tf) | |
| return {"ok": True, "items": df.to_dict(orient="records")} | |
| def get_config(): | |
| return {"ok": True, "config": RUNTIME_CONFIG} | |
| async def set_config(request: Request): | |
| try: | |
| payload = await request.json() | |
| except Exception: | |
| raise HTTPException(status_code=400, detail="Invalid JSON") | |
| if not isinstance(payload, dict): | |
| raise HTTPException(status_code=400, detail="JSON object expected") | |
| if "fusion" in payload and isinstance(payload["fusion"], dict): | |
| RUNTIME_CONFIG["fusion"].update(payload["fusion"]) | |
| if "weights" in payload and isinstance(payload["weights"], dict): | |
| new_w = {} | |
| for k, v in payload["weights"].items(): | |
| try: | |
| new_w[str(k)] = float(v) | |
| except Exception: | |
| continue | |
| if new_w: | |
| RUNTIME_CONFIG["weights"].update(new_w) | |
| return {"ok": True, "config": RUNTIME_CONFIG} | |
| async def perf_update(request: Request): | |
| try: | |
| p = await request.json() | |
| except Exception: | |
| raise HTTPException(status_code=400, detail="Invalid JSON") | |
| src = str(p.get("source") or "unknown") | |
| sig = str(p.get("signal") or "") | |
| move = float(p.get("actual_price_move") or 0) | |
| source_perf.update_performance(src, sig, move) | |
| return {"ok": True, "source": src, "score": source_perf.source_scores[src]} | |
| def perf_scores(): | |
| df = source_perf.as_dataframe() | |
| return {"ok": True, "items": df.to_dict(orient="records"), "overall": source_perf.overall_metrics()} | |
| def market_context(ticker: str = Query("BTCUSD"), timeframe: str = Query("5")): | |
| info = regime_detector.detect_from_alerts(ticker, timeframe) | |
| sess = session_label_utc() | |
| recent = [r for r in ALERTS if r.get("ticker") == ticker and str(r.get("tf")) == str(timeframe)][-30:] | |
| avg_strength = int(np.mean([r.get("strength") or 0 for r in recent])) if recent else 0 | |
| trigger, thr = optimize_alerts_threshold(avg_strength, info["regime"], sess) | |
| return {"ok": True, "ticker": ticker, "tf": str(timeframe), "session": sess, | |
| "regime": info, "avg_strength": avg_strength, "dynamic_threshold": thr, "trigger_now": trigger} | |
| async def ws_endpoint(ws: WebSocket): | |
| await ws_manager.connect(ws) | |
| try: | |
| while True: | |
| await ws.receive_text() | |
| except WebSocketDisconnect: | |
| ws_manager.disconnect(ws) | |
| # ========================= | |
| # GRADIO UI | |
| # ========================= | |
| example_ultra = """{ | |
| "source":"tradingview","emitter":"ULTRA-X-5m-ENH","ticker":"BTCUSD","tf":"5", | |
| "time":"2025-11-12T13:20:00Z","price":45000,"rsi":45.05,"adx_val":23.14, | |
| "vol_ratio":0.915,"macd_val":-5.5954,"macd_sig":40.3689,"trend_score":-2, | |
| "is_trending":true,"long_condition":false,"short_condition":false | |
| }""" | |
| def ui_refresh(ticker_filter, limit): | |
| try: | |
| df = latest_table(ticker_filter if ticker_filter else None, int(limit)) | |
| except Exception: | |
| df = pd.DataFrame() | |
| return df | |
| def ui_stats_md(): | |
| s = compute_stats() | |
| return (f"**📊 Σύνολο δεικτών:** {s['unique_emitters']} • " | |
| f"**📈 Σύνολο tickers:** {s['unique_tickers']} • " | |
| f"**🔔 Σύνολο σημάτων:** {s['total_events']} • " | |
| f"🟣 **Ενεργοί (30')**: {s['active_30']} • " | |
| f"🟪 **Ενεργοί (5')**: {s['active_5']} • " | |
| f"🕒 Τελευταία ενημέρωση (Athens): {datetime.now(tz=ATHENS_TZ).strftime('%H:%M:%S')}") | |
| def ui_stats_table(): | |
| s = compute_stats() | |
| return pd.DataFrame(s["per_emitter"]) | |
| def ui_universe(minutes, ticker): | |
| return active_sources(minutes=int(minutes), ticker=(ticker or None)) | |
| def ui_export_link(ticker, limit): | |
| limit = int(limit) | |
| ticker_q = (ticker or "").strip() | |
| href = f"/api/export.csv?limit={limit}&ticker={ticker_q}" | |
| return ( | |
| f"<a href='{href}' target='_blank' " | |
| f"style='display:inline-block;padding:10px 16px;border-radius:8px;" | |
| f"background:#6c63ff;color:white;text-decoration:none;font-weight:600;'>" | |
| f"📥 Download CSV</a>" | |
| ) | |
| def _render_fusion_card(fusion_result: Dict[str, Any], ticker: str, timeframe: str, icon: str) -> str: | |
| signal = fusion_result["fusion_signal"] | |
| signal_color = { | |
| "STRONG LONG": "#4CAF50", | |
| "MODERATE LONG": "#8BC34A", | |
| "NEUTRAL": "#FFC107", | |
| "MODERATE SHORT": "#FF9800", | |
| "STRONG SHORT": "#f44336", | |
| "NO SIGNAL": "#9E9E9E", | |
| "WEAK SIGNAL": "#FF5722" | |
| } | |
| for k in list(signal_color.keys()): | |
| if signal.startswith("MEGA ") and k in signal: | |
| signal_color[signal] = signal_color[k] | |
| color = signal_color.get(signal, "#9E9E9E") | |
| confidence_percent = fusion_result.get("fusion_confidence", 0) * 100 | |
| consensus_percent = fusion_result.get("consensus", 0) * 100 | |
| extra_html = "" | |
| if "enhanced_confidence" in fusion_result: | |
| extra_html = f""" | |
| <div style="display:grid;grid-template-columns:repeat(3,1fr);gap:15px;margin-top:10px;"> | |
| <div style="text-align:center;background:#111;padding:12px;border-radius:10px;"> | |
| <div style="font-size:22px;font-weight:bold;color:{color};">{fusion_result.get('enhanced_confidence',0):.2f}</div> | |
| <div style="font-size:12px;color:#aaa;">Enhanced Conf.</div> | |
| </div> | |
| <div style="text-align:center;background:#111;padding:12px;border-radius:10px;"> | |
| <div style="font-size:22px;font-weight:bold;color:{color};">{fusion_result.get('quality_score',0):.2f}</div> | |
| <div style="font-size:12px;color:#aaa;">Quality</div> | |
| </div> | |
| <div style="text-align:center;background:#111;padding:12px;border-radius:10px;"> | |
| <div style="font-size:22px;font-weight:bold;color:{color};">{fusion_result.get('composite_score',0):.2f}</div> | |
| <div style="font-size:12px;color:#aaa;">Composite</div> | |
| </div> | |
| </div> | |
| <div style="margin-top:10px;color:#ddd;"> | |
| Regime: <b>{fusion_result.get('market_regime')}</b> • Session: <b>{fusion_result.get('session')}</b> • | |
| Thr: <b>{fusion_result.get('dynamic_threshold')}</b> • Trigger: <b>{'✅' if fusion_result.get('trigger') else '❌'}</b> | |
| </div> | |
| """ | |
| html = f""" | |
| <div style="padding:25px;background:linear-gradient(135deg,{color}20,{color}40); | |
| border:2px solid {color};border-radius:15px;margin:15px 0;"> | |
| <div style="text-align:center;"> | |
| <h2 style="margin:0;color:{color};font-size:28px;">{icon} {signal}</h2> | |
| <div style="font-size:16px;color:#ccc;margin:10px 0;">{ticker} | {timeframe}</div> | |
| </div> | |
| <div style="display:grid;grid-template-columns:repeat(3,1fr);gap:15px;margin:20px 0;"> | |
| <div style="text-align:center;background:#111;padding:15px;border-radius:10px;"> | |
| <div style="font-size:24px;font-weight:bold;color:{color};">{confidence_percent:.1f}%</div> | |
| <div style="font-size:12px;color:#aaa;">Εμπιστοσύνη</div> | |
| </div> | |
| <div style="text-align:center;background:#111;padding:15px;border-radius:10px;"> | |
| <div style="font-size:24px;font-weight:bold;color:{color};">{consensus_percent:.1f}%</div> | |
| <div style="font-size:12px;color:#aaa;">Ομοφωνία</div> | |
| </div> | |
| <div style="text-align:center;background:#111;padding:15px;border-radius:10px;"> | |
| <div style="font-size:24px;font-weight:bold;color:{color};">{fusion_result.get('sources_used',0)}</div> | |
| <div style="font-size:12px;color:#aaa;">Πηγές</div> | |
| </div> | |
| </div> | |
| <div style="background:#111;padding:15px;border-radius:10px;margin:15px 0;"> | |
| <div style="font-weight:bold;margin-bottom:10px;">📋 Σύσταση:</div> | |
| <div style="color:#ddd;">{fusion_result.get('recommendation','')}</div> | |
| </div> | |
| <div style="background:#111;padding:15px;border-radius:10px;"> | |
| <div style="font-weight:bold;margin-bottom:10px;">🔍 Σύνθεση σήματος:</div> | |
| """ | |
| for source in fusion_result.get("source_details", []): | |
| contrib_color = "#4CAF50" if source["contribution"] > 0 else "#f44336" if source["contribution"] < 0 else "#FFC107" | |
| html += f""" | |
| <div style="display:flex;justify-content:space-between;margin:5px 0;padding:5px;background:#0f0f0f;border-radius:5px;"> | |
| <span>{source['source']}</span> | |
| <span style="color:{contrib_color};font-weight:bold;">{source['signal']} ({source['strength']}%)</span> | |
| <span>βάρος: {source['weight']}</span> | |
| </div> | |
| """ | |
| html += f"</div>{extra_html}</div>" | |
| return html | |
| def ui_fusion_signal(ticker, timeframe, mode: str = "Enhanced"): | |
| mode = (mode or "Enhanced") | |
| if mode == "Mega": | |
| res = mega_engine.process_mega_fusion(ticker, timeframe) | |
| log_fusion_result(res, "Mega") | |
| return _render_fusion_card(res, ticker, timeframe, "🦖") | |
| if mode == "Pro+": | |
| res = pro_plus_engine.calculate_pro_plus_fusion(ticker, timeframe) | |
| log_fusion_result(res, "Pro+") | |
| return _render_fusion_card(res, ticker, timeframe, "🧠") | |
| if mode == "Enhanced": | |
| res = enhanced_engine.calculate_enhanced_fusion(ticker, timeframe) | |
| log_fusion_result(res, "Enhanced") | |
| return _render_fusion_card(res, ticker, timeframe, "🚀") | |
| if mode == "Base": | |
| res = enhanced_engine.calculate_fusion_signal(ticker, timeframe) | |
| log_fusion_result(res, "Base") | |
| return _render_fusion_card(res, ticker, timeframe, "🎯") | |
| # default fallback | |
| res = enhanced_engine.calculate_enhanced_fusion(ticker, timeframe) | |
| log_fusion_result(res, "Enhanced") | |
| return _render_fusion_card(res, ticker, timeframe, "🚀") | |
| def analyze_json(json_text: str) -> str: | |
| try: | |
| data = json.loads(json_text) | |
| except Exception as e: | |
| return f"❌ Invalid JSON: {e}" | |
| norm = normalize_payload(data) | |
| pretty = json.dumps({k: v for k, v in norm.items() if k != "raw"}, ensure_ascii=False, indent=2) | |
| return f"**Schema:** `{norm['schema']}` \n**Signal:** `{norm['signal']}` \n**Strength:** `{norm['strength']}` \n\n```json\n{pretty}\n```" | |
| def ui_save_settings(tw, minc, wthr, cthr, weights_json): | |
| RUNTIME_CONFIG["fusion"].update({ | |
| "time_window_minutes": int(tw), | |
| "min_confidence": float(minc), | |
| "weight_threshold": float(wthr), | |
| "consensus_threshold": float(cthr), | |
| }) | |
| try: | |
| w = json.loads(weights_json) if weights_json.strip() else {} | |
| new_w = {} | |
| for k, v in w.items(): | |
| try: | |
| new_w[str(k)] = float(v) | |
| except Exception: | |
| continue | |
| if new_w: | |
| RUNTIME_CONFIG["weights"].update(new_w) | |
| msg = "✅ Settings saved." | |
| except Exception as e: | |
| msg = f"⚠️ Weights JSON error: {e}" | |
| return msg, json.dumps(RUNTIME_CONFIG, ensure_ascii=False, indent=2) | |
| def ui_perf_table(): | |
| return source_perf.as_dataframe() | |
| def ui_perf_summary_md(): | |
| m = source_perf.overall_metrics() | |
| return f"**Signals:** {m['total_signals']} • **Wins:** {m['wins']} • **Win‑Rate:** {m['win_rate']}%" | |
| def ui_fusion_history(limit: int): | |
| return fusion_history_df(limit) | |
| # ------------------------- | |
| # UI Build | |
| # ------------------------- | |
| with gr.Blocks( | |
| title="ULTRA‑X Multi‑Source Signal Board Pro", | |
| theme=gr.themes.Soft(), | |
| css=""" | |
| body { background:#0b0c10; color:#eaeaea; } | |
| .dashboard-stats { background: linear-gradient(135deg, #667eea 0%, #764ba2 100%); color: white; | |
| padding: 20px; border-radius: 10px; margin: 10px 0; } | |
| """ | |
| ) as demo: | |
| gr.Markdown("# 🚀 ULTRA‑X Multi‑Source Signal Board Pro") | |
| gr.Markdown("Real-time trading signals from **ULTRA‑X Feeder, Master BTC, QSC Pro, RTE, BOS**, and more") | |
| # ---------- Fusion Dashboard ---------- | |
| with gr.Tab("🎯 Fusion Dashboard"): | |
| gr.Markdown("### Συγκεντρωτικό Σήμα Trading") | |
| with gr.Row(): | |
| with gr.Column(scale=1): | |
| fusion_ticker = gr.Textbox(label="Ticker", value="BTCUSD") | |
| fusion_timeframe = gr.Dropdown(label="Timeframe", choices=["1", "5", "15", "60", "240", "1D"], value="5") | |
| fusion_mode = gr.Dropdown( | |
| label="Mode", | |
| choices=["Mega", "Pro+", "Enhanced", "Base"], | |
| value="Mega" | |
| ) | |
| fusion_refresh = gr.Button("🔄 Ανανέωση Fusion", variant="primary") | |
| with gr.Column(scale=2): | |
| fusion_display = gr.HTML(value=ui_fusion_signal("BTCUSD", "5", "Mega")) | |
| fusion_refresh.click(fn=ui_fusion_signal, | |
| inputs=[fusion_ticker, fusion_timeframe, fusion_mode], | |
| outputs=fusion_display) | |
| demo.load(fn=ui_fusion_signal, | |
| inputs=[fusion_ticker, fusion_timeframe, fusion_mode], | |
| outputs=fusion_display, | |
| every=30) | |
| gr.Markdown("### 📜 Fusion History (Athens Time)") | |
| hist_size = gr.Slider(10, 200, value=50, step=5, label="History size") | |
| hist_table = gr.Dataframe(headers=["when", "mode", "ticker", "tf", "signal", "conf%", "dyn_conf%", "sources", "regime", "session"], | |
| wrap=True, height=260, label="Fusion History") | |
| hist_size.change(fn=ui_fusion_history, inputs=hist_size, outputs=hist_table) | |
| demo.load(fn=ui_fusion_history, inputs=hist_size, outputs=hist_table, every=20) | |
| # ---------- Live Dashboard ---------- | |
| with gr.Tab("📡 Live Dashboard"): | |
| stats_md = gr.Markdown(elem_classes="dashboard-stats") | |
| with gr.Row(): | |
| with gr.Column(scale=2): | |
| with gr.Row(): | |
| t = gr.Textbox(label="🔍 Filter Ticker", value="BTCUSD", container=False) | |
| l = gr.Slider(minimum=50, maximum=1000, value=200, step=50, label="📈 Display Limit") | |
| with gr.Column(scale=1): | |
| refresh_btn = gr.Button("🔄 Refresh Now", variant="primary") | |
| export_btn = gr.Button("📥 Export CSV", variant="secondary") | |
| grid = gr.Dataframe( | |
| label="🎯 Live Trading Signals", | |
| headers=["when", "src", "ticker", "tf", "signal", "strength", "price", "SL", "TP", "RR", "pos%", "ADX", "RSI", "bar_time"], | |
| wrap=True, | |
| height=420 | |
| ) | |
| stats_df = gr.Dataframe( | |
| label="📊 Signals by Source", | |
| headers=["src", "signals", "tickers", "tfs", "last_seen", "last_bar_time"], | |
| wrap=True | |
| ) | |
| with gr.Accordion("🌌 Active Universe (sources per pair)", open=True): | |
| uni_minutes = gr.Slider(5, 240, value=60, step=5, label="Active Universe Window (min)") | |
| uni_table = gr.Dataframe(headers=["ticker", "tf", "sources_count", "sources", "last_seen"], wrap=True) | |
| refresh_btn.click(fn=ui_refresh, inputs=[t, l], outputs=grid) | |
| refresh_btn.click(fn=ui_stats_md, inputs=None, outputs=stats_md) | |
| refresh_btn.click(fn=ui_stats_table, inputs=None, outputs=stats_df) | |
| refresh_btn.click(fn=ui_universe, inputs=[uni_minutes, t], outputs=uni_table) | |
| demo.load(fn=ui_refresh, inputs=[t, l], outputs=grid, every=10) | |
| demo.load(fn=ui_stats_md, inputs=None, outputs=stats_md, every=10) | |
| demo.load(fn=ui_stats_table, inputs=None, outputs=stats_df, every=10) | |
| demo.load(fn=ui_universe, inputs=[uni_minutes, t], outputs=uni_table, every=10) | |
| download_html = gr.HTML() | |
| export_btn.click(fn=ui_export_link, inputs=[t, l], outputs=download_html) | |
| # ---------- Enhanced Analytics ---------- | |
| with gr.Tab("📊 Enhanced Analytics"): | |
| gr.Markdown("### Απόδοση Πηγών (Dynamic Weights)") | |
| perf_md = gr.Markdown() | |
| perf_table = gr.Dataframe(headers=["source", "total_signals", "successful_signals", "win_rate_%", "current_weight"], | |
| wrap=True, height=300) | |
| refresh_perf = gr.Button("🔄 Refresh Performance") | |
| refresh_perf.click(fn=ui_perf_summary_md, inputs=None, outputs=perf_md) | |
| refresh_perf.click(fn=ui_perf_table, inputs=None, outputs=perf_table) | |
| demo.load(fn=ui_perf_summary_md, inputs=None, outputs=perf_md, every=15) | |
| demo.load(fn=ui_perf_table, inputs=None, outputs=perf_table, every=15) | |
| gr.Markdown("### Market Context (Regime / Threshold)") | |
| mc_t = gr.Textbox(label="Ticker", value="BTCUSD") | |
| mc_tf = gr.Dropdown(label="Timeframe", choices=["1", "5", "15", "60", "240", "1D"], value="5") | |
| mc_btn = gr.Button("🔎 Get Market Context", variant="secondary") | |
| mc_box = gr.Code(label="Context JSON (readonly)") | |
| def ui_market_context(ticker, tf): | |
| ctx = market_context(ticker, tf) | |
| return json.dumps(ctx, ensure_ascii=False, indent=2) | |
| mc_btn.click(fn=ui_market_context, inputs=[mc_t, mc_tf], outputs=mc_box) | |
| # ---------- Signal Analyzer ---------- | |
| with gr.Tab("🧪 Signal Analyzer"): | |
| gr.Markdown("### Test & Analyze JSON Payloads") | |
| with gr.Row(): | |
| with gr.Column(scale=1): | |
| ex = gr.Dropdown(choices=["ULTRA"], value="ULTRA", label="Sample Template") | |
| generate_btn = gr.Button("🎲 Load Sample", variant="primary") | |
| with gr.Row(): | |
| inp = gr.Textbox(label="JSON Payload", lines=12, value=example_ultra, show_copy_button=True) | |
| with gr.Row(): | |
| run = gr.Button("🚀 Analyze & Normalize", variant="primary") | |
| out = gr.Markdown(label="Analysis Results") | |
| def set_example(kind): | |
| return example_ultra | |
| ex.change(fn=set_example, inputs=ex, outputs=inp) | |
| generate_btn.click(fn=set_example, inputs=ex, outputs=inp) | |
| run.click(fn=analyze_json, inputs=inp, outputs=out) | |
| # ---------- Fusion Settings ---------- | |
| with gr.Tab("🛠️ Fusion Settings"): | |
| gr.Markdown("### Ρυθμίσεις Fusion Engine (χωρίς redeploy)") | |
| f = RUNTIME_CONFIG["fusion"] | |
| w = RUNTIME_CONFIG["weights"] | |
| tw = gr.Slider(1, 60, value=f["time_window_minutes"], step=1, label="Time Window (min)") | |
| minc = gr.Slider(0.1, 1.0, value=f["min_confidence"], step=0.05, label="Min Confidence") | |
| wthr = gr.Slider(0.0, 1.0, value=f["weight_threshold"], step=0.01, label="Weight Threshold") | |
| cthr = gr.Slider(0.1, 1.0, value=f["consensus_threshold"], step=0.05, label="Consensus Threshold") | |
| weights_json = gr.Code(language="json", value=json.dumps(w, indent=2), label="Source Weights (JSON)") | |
| save_btn = gr.Button("💾 Save Settings", variant="primary") | |
| save_msg = gr.Markdown() | |
| curr_cfg = gr.Code(label="Current Settings", value=json.dumps(RUNTIME_CONFIG, ensure_ascii=False, indent=2)) | |
| save_btn.click(fn=ui_save_settings, inputs=[tw, minc, wthr, cthr, weights_json], outputs=[save_msg, curr_cfg]) | |
| # ---------- Configuration ---------- | |
| with gr.Tab("⚙️ Configuration"): | |
| gr.Markdown("### Setup Guide & API") | |
| gr.Markdown( | |
| """ | |
| **Webhook:** `POST /api/webhook?token=YOUR_TOKEN` | |
| **Latest:** `GET /api/latest` | |
| **Fusion (base):** `GET /api/fusion?ticker=BTCUSD&timeframe=5` | |
| **Fusion (enhanced):** `GET /api/fusion/enhanced?ticker=BTCUSD&timeframe=5` | |
| **Fusion (Pro+):** `GET /api/fusion/proplus?ticker=BTCUSD&timeframe=5` | |
| **Fusion (Mega):** `GET /api/fusion/mega?ticker=BTCUSD&timeframe=5` | |
| **Active Sources:** `GET /api/sources/active?minutes=5&ticker=BTCUSD&tf=5` | |
| **Config (runtime):** `GET/POST /api/config` | |
| **Performance:** `POST /api/perf/update`, `GET /api/perf/scores` | |
| **Market Context:** `GET /api/market/context?ticker=BTCUSD&timeframe=5` | |
| """ | |
| ) | |
| # Mount Gradio app στο "/" | |
| app = gr.mount_gradio_app(api, demo, path="/") | |
| async def on_start(): | |
| print("🚀 ULTRA‑X Multi‑Source Signal Board Pro (Mega Edition) starting…") | |
| print(f"🔐 WEBHOOK_TOKEN set: {bool(WEBHOOK_TOKEN)}") | |
| print(f"💾 Persist path: {PERSIST_PATH or 'disabled'} (max {MAX_LOG_MB}MB)") | |