CoolFace
Apppublic

shaibu01/Titan-Engine

sourceHugging Faceupdated 2d agoView on Hugging Face
0likes
dynamic_alpha_hunter.py500 linesDownload Raw Back to root
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