trenderbot / app.py
michalis13's picture
Update app.py
de7ac92 verified
Raw
History Blame Contribute Delete
82.5 kB
# 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})
@staticmethod
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"])
}
@staticmethod
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
@staticmethod
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)
@staticmethod
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": "🔍 Δεν υπάρχουν πρόσφατα σήματα"
}
@staticmethod
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
@staticmethod
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))
@staticmethod
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=["*"])
@api.get("/api/healthz")
def healthz():
return {"ok": True, "alerts_cached": len(ALERTS), "persist": bool(PERSIST_PATH)}
@api.get("/api/")
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")
@api.post("/api/webhook")
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")
@api.get("/api/latest")
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")}
@api.get("/api/export.csv")
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"})
@api.get("/api/stats")
def stats():
s = compute_stats()
return {"ok": True, **s}
@api.get("/api/fusion")
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}
@api.get("/api/fusion/enhanced")
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}
@api.get("/api/fusion/proplus")
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}
@api.get("/api/fusion/mega")
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}
@api.get("/api/sources/active")
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")}
@api.get("/api/config")
def get_config():
return {"ok": True, "config": RUNTIME_CONFIG}
@api.post("/api/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}
@api.post("/api/perf/update")
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]}
@api.get("/api/perf/scores")
def perf_scores():
df = source_perf.as_dataframe()
return {"ok": True, "items": df.to_dict(orient="records"), "overall": source_perf.overall_metrics()}
@api.get("/api/market/context")
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}
@api.websocket("/ws")
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="/")
@api.on_event("startup")
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)")