shaibu01/Titan-Engine
0
1"""2Provider-neutral dynamic alpha hunter for Kraken, Solana/Jupiter, and event markets.3 4The module is deliberately execution-safe: it scans, records anomalies, builds5execution blueprints, and publishes redacted state. Actual broker/DEX/order6submission remains behind the existing provider-specific execution gates.7"""8 9from __future__ import annotations10 11import asyncio12import math13import os14import threading15import time16from collections import deque17from dataclasses import asdict, dataclass, field18from typing import Any, Deque, Dict, Iterable, List, Optional19 20try:21 import aiohttp22except Exception: # pragma: no cover - optional in stripped runtimes23 aiohttp = None24 25try:26 from polymarket_oracle import polymarket_market_relevance27except Exception: # pragma: no cover - keep hunter importable in partial deploys28 def polymarket_market_relevance(market: Dict[str, Any]) -> Dict[str, Any]:29 return {30 "eligible": True,31 "reason": "RELEVANCE_FILTER_UNAVAILABLE",32 "relevance_score": 0.20,33 "cluster_key": str((market or {}).get("slug") or (market or {}).get("question") or "POLYMARKET"),34 }35 36 37def _env_bool(name: str, default: bool = False) -> bool:38 raw = str(os.getenv(name, "1" if default else "0") or "").strip().lower()39 return raw in {"1", "true", "yes", "on", "live", "enabled"}40 41 42def _safe_float(value: Any, default: float = 0.0) -> float:43 try:44 if value is None:45 return default46 result = float(value)47 if math.isfinite(result):48 return result49 except Exception:50 pass51 return default52 53 54def _csv_env(name: str) -> List[str]:55 return [item.strip() for item in str(os.getenv(name, "") or "").split(",") if item.strip()]56 57 58@dataclass59class AnomalySignal:60 scanner: str61 kind: str62 asset: str63 venue: str = ""64 reference_venue: str = ""65 score: float = 0.066 route: str = "RESEARCH"67 reason: str = ""68 metrics: Dict[str, Any] = field(default_factory=dict)69 ts: float = field(default_factory=time.time)70 71 72@dataclass73class ScannerState:74 name: str75 status: str = "BOOTING"76 targets: int = 077 anomalies: int = 078 last_error: str = ""79 last_update: float = 0.080 81 82@dataclass83class ExecutionPlan:84 accepted: bool85 route: str86 asset: str87 side: str = "WATCH"88 reason: str = ""89 max_slippage_bps: float = 75.090 volatility_stop_atr_mult: float = 1.591 liquidity_decay_depth_floor: float = 0.4092 time_kill_switch_sec: int = 90093 metadata: Dict[str, Any] = field(default_factory=dict)94 95 96class SystemSharedMemory:97 """Small thread-safe shared memory surface consumed by cloud desks."""98 99 def __init__(self, max_items: int = 200):100 self._lock = threading.RLock()101 self.targets: Deque[Dict[str, Any]] = deque(maxlen=max_items)102 self.anomalies: Deque[Dict[str, Any]] = deque(maxlen=max_items)103 self.positions: Deque[Dict[str, Any]] = deque(maxlen=max_items)104 self.volatility_metrics: Deque[Dict[str, Any]] = deque(maxlen=max_items)105 106 def record_target(self, payload: Dict[str, Any]) -> None:107 with self._lock:108 self.targets.append({"ts": time.time(), **dict(payload or {})})109 110 def record_anomaly(self, payload: Dict[str, Any]) -> None:111 with self._lock:112 self.anomalies.append({"ts": time.time(), **dict(payload or {})})113 114 def record_position(self, payload: Dict[str, Any]) -> None:115 with self._lock:116 self.positions.append({"ts": time.time(), **dict(payload or {})})117 118 def record_volatility_metric(self, payload: Dict[str, Any]) -> None:119 with self._lock:120 self.volatility_metrics.append({"ts": time.time(), **dict(payload or {})})121 122 def snapshot(self) -> Dict[str, Any]:123 with self._lock:124 return {125 "targets": list(self.targets)[-30:],126 "anomalies": list(self.anomalies)[-30:],127 "positions": list(self.positions)[-30:],128 "volatility_metrics": list(self.volatility_metrics)[-30:],129 }130 131 132class HeuristicEvaluator:133 """Converts anomalies into provider-specific execution blueprints."""134 135 def __init__(136 self,137 min_depth_usd: float = 50_000.0,138 min_score: float = 0.65,139 max_slippage_bps: float = 75.0,140 ):141 self.min_depth_usd = float(min_depth_usd)142 self.min_score = float(min_score)143 self.max_slippage_bps = float(max_slippage_bps)144 145 def evaluate(self, anomaly: AnomalySignal, market_depth: Optional[Dict[str, Any]] = None) -> ExecutionPlan:146 depth = _safe_float((market_depth or {}).get("depth_1pct_usd"), 0.0)147 slippage = _safe_float((market_depth or {}).get("slippage_bps"), 0.0)148 if depth and depth < self.min_depth_usd:149 return ExecutionPlan(150 accepted=False,151 route=anomaly.route,152 asset=anomaly.asset,153 reason=f"DEPTH_TOO_THIN:{depth:.0f}<{self.min_depth_usd:.0f}",154 metadata=asdict(anomaly),155 )156 if slippage and slippage > self.max_slippage_bps:157 return ExecutionPlan(158 accepted=False,159 route=anomaly.route,160 asset=anomaly.asset,161 reason=f"SLIPPAGE_TOO_HIGH:{slippage:.1f}>{self.max_slippage_bps:.1f}",162 metadata=asdict(anomaly),163 )164 if anomaly.score < self.min_score:165 return ExecutionPlan(166 accepted=False,167 route=anomaly.route,168 asset=anomaly.asset,169 reason=f"SCORE_TOO_LOW:{anomaly.score:.3f}<{self.min_score:.3f}",170 metadata=asdict(anomaly),171 )172 side = "BUY"173 if anomaly.kind in {"negative_isolation_spread", "funding_short", "event_overpriced"}:174 side = "SELL"175 return ExecutionPlan(176 accepted=True,177 route=anomaly.route,178 asset=anomaly.asset,179 side=side,180 reason="HEURISTIC_BLUEPRINT_ACCEPTED",181 max_slippage_bps=self.max_slippage_bps,182 metadata=asdict(anomaly),183 )184 185 186class DynamicAlphaHunter:187 """Async scanner coordinator with redacted runtime state."""188 189 def __init__(self, enabled: Optional[bool] = None):190 self.enabled = _env_bool("TITAN_DYNAMIC_ALPHA_HUNTER_ENABLED", True) if enabled is None else bool(enabled)191 self.scan_interval_sec = float(os.getenv("TITAN_DYNAMIC_ALPHA_SCAN_INTERVAL_SEC", "12"))192 self.min_pool_liquidity_usd = float(os.getenv("TITAN_SOLANA_MIN_POOL_LIQUIDITY_USD", "50000"))193 self.min_volume_z = float(os.getenv("TITAN_SOLANA_VOLUME_Z_MIN", "3.0"))194 self.max_anomalies = int(os.getenv("TITAN_DYNAMIC_ALPHA_MAX_ANOMALIES", "100"))195 self.memory = SystemSharedMemory(max_items=300)196 self.evaluator = HeuristicEvaluator(197 min_depth_usd=float(os.getenv("TITAN_DYNAMIC_ALPHA_MIN_DEPTH_USD", "50000")),198 min_score=float(os.getenv("TITAN_DYNAMIC_ALPHA_MIN_SCORE", "0.65")),199 max_slippage_bps=float(os.getenv("TITAN_DYNAMIC_ALPHA_MAX_SLIPPAGE_BPS", "75")),200 )201 self.recent_anomalies: Deque[Dict[str, Any]] = deque(maxlen=self.max_anomalies)202 self.recent_plans: Deque[Dict[str, Any]] = deque(maxlen=self.max_anomalies)203 self.scanners: Dict[str, ScannerState] = {204 "exchange_isolation": ScannerState(name="exchange_isolation"),205 "onchain_liquidity_velocity": ScannerState(name="onchain_liquidity_velocity"),206 "derivatives_squeeze": ScannerState(name="derivatives_squeeze"),207 "event_probability_discrepancy": ScannerState(name="event_probability_discrepancy"),208 }209 self._thread: Optional[threading.Thread] = None210 self._running = False211 self._lock = threading.RLock()212 213 @property214 def secret_status(self) -> Dict[str, Any]:215 return {216 "kraken_api_key_configured": bool(os.getenv("KRAKEN_API_KEY", "").strip()),217 "kraken_private_key_configured": bool(218 os.getenv("KRAKEN_PRIVATE_KEY", "").strip()219 or os.getenv("KRAKEN_API_SECRET", "").strip()220 ),221 "phantom_wallet_secret_name": "PHANTOM_SOLANA_WALLET_ADDRESS",222 "phantom_solana_wallet_configured": bool(223 os.getenv("PHANTOM_SOLANA_WALLET_ADDRESS", "").strip()224 or os.getenv("SOLANA_WALLET_ADDRESS", "").strip()225 ),226 "solana_rpc_configured": bool(os.getenv("SOLANA_RPC_URL", "").strip()),227 "polymarket_enabled": _env_bool("TITAN_POLY_ORACLE_ENABLED", True),228 }229 230 def start(self) -> "DynamicAlphaHunter":231 if not self.enabled:232 return self233 with self._lock:234 if self._running:235 return self236 self._running = True237 self._thread = threading.Thread(target=self._run_loop, name="dynamic_alpha_hunter", daemon=True)238 self._thread.start()239 return self240 241 def _run_loop(self) -> None:242 try:243 asyncio.run(self._main())244 except Exception as exc:245 with self._lock:246 for state in self.scanners.values():247 state.status = "ERROR"248 state.last_error = str(exc)[:240]249 state.last_update = time.time()250 self._running = False251 252 async def _main(self) -> None:253 await asyncio.gather(254 self._exchange_isolation_scanner(),255 self._onchain_liquidity_velocity_scanner(),256 self._derivatives_squeeze_scanner(),257 self._event_probability_scanner(),258 )259 260 async def _sleep(self, multiplier: float = 1.0) -> None:261 await asyncio.sleep(max(2.0, self.scan_interval_sec * multiplier))262 263 async def _get_json(self, url: str, timeout: float = 6.0) -> Optional[Any]:264 if aiohttp is None or not url:265 return None266 try:267 async with aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=timeout)) as session:268 async with session.get(url) as response:269 if response.status >= 400:270 return None271 return await response.json(content_type=None)272 except Exception:273 return None274 275 def _update_state(self, key: str, status: str, targets: int = 0, error: str = "") -> None:276 with self._lock:277 state = self.scanners[key]278 state.status = status279 state.targets = int(targets)280 state.last_error = str(error or "")[:240]281 state.last_update = time.time()282 283 def _record_anomaly(self, anomaly: AnomalySignal, depth: Optional[Dict[str, Any]] = None) -> None:284 plan = self.evaluator.evaluate(anomaly, depth)285 anomaly_payload = asdict(anomaly)286 plan_payload = asdict(plan)287 with self._lock:288 self.recent_anomalies.append(anomaly_payload)289 self.recent_plans.append(plan_payload)290 if anomaly.scanner in self.scanners:291 self.scanners[anomaly.scanner].anomalies += 1292 self.scanners[anomaly.scanner].last_update = time.time()293 self.memory.record_anomaly(anomaly_payload)294 self.memory.record_target(plan_payload)295 296 async def _exchange_isolation_scanner(self) -> None:297 keywords = ("suspend", "wallet maintenance", "deposit pause", "withdrawal pause", "network upgrade")298 urls = _csv_env("TITAN_EXCHANGE_ANNOUNCEMENT_URLS")299 while True:300 try:301 if not urls:302 self._update_state("exchange_isolation", "CONFIGURE_ANNOUNCEMENT_URLS", 0)303 await self._sleep(5)304 continue305 hits = 0306 for url in urls[:8]:307 payload = await self._get_json(url)308 text = str(payload or "").lower()309 if any(word in text for word in keywords):310 hits += 1311 self._record_anomaly(AnomalySignal(312 scanner="exchange_isolation",313 kind="exchange_isolation_spread",314 asset="WATCHLIST",315 venue=url,316 score=0.72,317 route="KRAKEN",318 reason="announcement_keyword_match",319 metrics={"keywords": [word for word in keywords if word in text]},320 ))321 self._update_state("exchange_isolation", "WATCHING", len(urls))322 except Exception as exc:323 self._update_state("exchange_isolation", "ERROR", len(urls), str(exc))324 await self._sleep(4)325 326 async def _onchain_liquidity_velocity_scanner(self) -> None:327 tokens = _csv_env("TITAN_SOLANA_TOKEN_WATCHLIST") or ["SOL", "JUP", "WIF", "BONK"]328 dexscreener_base = os.getenv("DEXSCREENER_API_BASE_URL", "https://api.dexscreener.com/latest/dex/search?q=").rstrip("=")329 while True:330 try:331 if aiohttp is None:332 self._update_state("onchain_liquidity_velocity", "AIOHTTP_UNAVAILABLE", len(tokens))333 await self._sleep(5)334 continue335 for token in tokens[:12]:336 url = f"{dexscreener_base}={token}"337 payload = await self._get_json(url)338 pairs = payload.get("pairs", []) if isinstance(payload, dict) else []339 best = max(340 (p for p in pairs if isinstance(p, dict)),341 key=lambda p: _safe_float((p.get("liquidity") or {}).get("usd"), 0.0),342 default=None,343 )344 if not best:345 continue346 liquidity = _safe_float((best.get("liquidity") or {}).get("usd"), 0.0)347 volume_5m = _safe_float((best.get("volume") or {}).get("m5"), 0.0)348 volume_1h = _safe_float((best.get("volume") or {}).get("h1"), 0.0)349 baseline = max(volume_1h / 12.0, 1.0)350 z_vol = (volume_5m - baseline) / max(math.sqrt(baseline), 1.0)351 self.memory.record_volatility_metric({352 "provider": "SOLANA_DEX",353 "asset": token,354 "liquidity_usd": liquidity,355 "volume_5m": volume_5m,356 "volume_z": z_vol,357 })358 if liquidity >= self.min_pool_liquidity_usd and z_vol >= self.min_volume_z:359 self._record_anomaly(360 AnomalySignal(361 scanner="onchain_liquidity_velocity",362 kind="liquidity_velocity",363 asset=str(best.get("baseToken", {}).get("symbol") or token),364 venue=str(best.get("dexId") or "DEXSCREENER"),365 score=min(1.0, 0.55 + (z_vol / 10.0)),366 route="SOLANA_DEX",367 reason="5m_volume_velocity_spike",368 metrics={"liquidity_usd": liquidity, "volume_z": z_vol},369 ),370 {"depth_1pct_usd": liquidity, "slippage_bps": 75.0},371 )372 self._update_state("onchain_liquidity_velocity", "WATCHING_SOLANA_POOLS", len(tokens))373 except Exception as exc:374 self._update_state("onchain_liquidity_velocity", "ERROR", len(tokens), str(exc))375 await self._sleep(1)376 377 async def _derivatives_squeeze_scanner(self) -> None:378 symbols = _csv_env("TITAN_DERIVATIVES_SQUEEZE_WATCHLIST") or ["BTC/USD", "ETH/USD", "SOL/USD"]379 while True:380 try:381 status = "KRAKEN_KEYS_READY" if self.secret_status["kraken_api_key_configured"] else "KRAKEN_PUBLIC_WATCH"382 self._update_state("derivatives_squeeze", status, len(symbols))383 except Exception as exc:384 self._update_state("derivatives_squeeze", "ERROR", len(symbols), str(exc))385 await self._sleep(2)386 387 async def _event_probability_scanner(self) -> None:388 markets_url = os.getenv("POLYMARKET_GAMMA_MARKETS_URL", "https://gamma-api.polymarket.com/markets")389 while True:390 try:391 if not _env_bool("TITAN_POLY_ORACLE_ENABLED", True):392 self._update_state("event_probability_discrepancy", "DISABLED", 0)393 await self._sleep(5)394 continue395 payload = await self._get_json(markets_url)396 markets: Iterable[Any]397 if isinstance(payload, list):398 markets = payload399 elif isinstance(payload, dict):400 markets = payload.get("markets", []) or payload.get("data", []) or []401 else:402 markets = []403 count = 0404 eligible_count = 0405 filtered_count = 0406 anomaly_count = 0407 seen_clusters = set()408 min_pressure = _safe_float(os.getenv("TITAN_POLY_EVENT_PRESSURE_RATIO_MIN", "4.0"), 4.0)409 for market in list(markets)[:50]:410 if not isinstance(market, dict):411 continue412 count += 1413 relevance = polymarket_market_relevance(market)414 if not relevance.get("eligible", False):415 filtered_count += 1416 continue417 cluster_key = str(relevance.get("cluster_key") or market.get("slug") or market.get("question") or count)418 if cluster_key in seen_clusters:419 filtered_count += 1420 continue421 seen_clusters.add(cluster_key)422 eligible_count += 1423 volume = _safe_float(market.get("volume24hr", market.get("volume24h", market.get("volume"))), 0.0)424 liquidity = _safe_float(market.get("liquidity"), 0.0)425 pressure = volume / max(liquidity, 1.0) if volume > 0 and liquidity > 0 else 0.0426 if pressure > min_pressure:427 relevance_score = _safe_float(relevance.get("relevance_score"), 0.0)428 anomaly_score = min(0.95, 0.46 + min(pressure, 12.0) / 24.0 + relevance_score * 0.08)429 anomaly_count += 1430 self._record_anomaly(AnomalySignal(431 scanner="event_probability_discrepancy",432 kind="event_probability_gap",433 asset=str(market.get("slug") or market.get("question") or "POLYMARKET"),434 venue="POLYMARKET",435 score=float(anomaly_score),436 route="POLYMARKET",437 reason="volume_to_liquidity_pressure",438 metrics={439 "volume": volume,440 "liquidity": liquidity,441 "pressure": pressure,442 "relevance": relevance,443 },444 ))445 self.memory.record_volatility_metric({446 "provider": "POLYMARKET",447 "scanner": "event_probability_discrepancy",448 "markets_scanned": count,449 "eligible_markets": eligible_count,450 "filtered_markets": filtered_count,451 "anomalies": anomaly_count,452 })453 self._update_state("event_probability_discrepancy", "WATCHING_RELEVANT_EVENT_MARKETS", eligible_count)454 except Exception as exc:455 self._update_state("event_probability_discrepancy", "ERROR", 0, str(exc))456 await self._sleep(3)457 458 def snapshot(self) -> Dict[str, Any]:459 with self._lock:460 return {461 "enabled": bool(self.enabled),462 "running": bool(self._running),463 "secret_status": self.secret_status,464 "scanner_states": {key: asdict(value) for key, value in self.scanners.items()},465 "recent_anomalies": list(self.recent_anomalies)[-25:],466 "recent_plans": list(self.recent_plans)[-25:],467 "shared_memory": self.memory.snapshot(),468 "routes": {469 "exchange_isolation": "KRAKEN",470 "onchain_liquidity_velocity": "SOLANA_DEX",471 "derivatives_squeeze": "KRAKEN",472 "event_probability_discrepancy": "POLYMARKET",473 },474 "config": {475 "scan_interval_sec": self.scan_interval_sec,476 "min_pool_liquidity_usd": self.min_pool_liquidity_usd,477 "min_volume_z": self.min_volume_z,478 "max_slippage_bps": self.evaluator.max_slippage_bps,479 },480 }481 482 483_HUNTER: Optional[DynamicAlphaHunter] = None484_HUNTER_LOCK = threading.RLock()485 486 487def start_dynamic_alpha_hunter() -> DynamicAlphaHunter:488 global _HUNTER489 with _HUNTER_LOCK:490 if _HUNTER is None:491 _HUNTER = DynamicAlphaHunter()492 return _HUNTER.start()493 494 495def dynamic_alpha_snapshot() -> Dict[str, Any]:496 with _HUNTER_LOCK:497 if _HUNTER is None:498 return DynamicAlphaHunter(enabled=False).snapshot()499 return _HUNTER.snapshot()500 