"""ML-powered signal generator for HK weather prediction markets. Uses three-layer calibrated LightGBM models (raw → Platt → isotonic) combined with spatial features, typhoon model, and portfolio Kelly. """ import sys from pathlib import Path from datetime import datetime, timedelta from typing import Dict, Optional, Tuple, List import numpy as np import pandas as pd sys.path.insert(0, str(Path(__file__).parent.parent)) from ml.features import FeatureEngine from ml.model import ModelEnsemble, TARGET_DEFINITIONS from ml.spatial import SpatialWeatherClient from ml.typhoon import TyphoonModel from weather.openmeteo_client import OpenMeteoClient from weather.hko_client import HKOClient from strategy.calibrator import ProbabilityCalibrator from strategy.kelly import KellyCriterion from strategy.portfolio_kelly import PortfolioKelly from config import HK_COORDS, MIN_EDGE_BPS, KELLY_FRACTION, MAX_POSITION_USDC class MLPredictor: """ ML-based weather probability predictor for Polymarket trading. Combines: 1. Feature engineering from NWP model output 2. Trained LightGBM probability models 3. Platt scaling calibration on historical outcomes 4. Ensemble disagreement as edge amplifier 5. Kelly criterion position sizing """ def __init__( self, bankroll_usdc: float = 1000.0, min_edge_bps: float = MIN_EDGE_BPS, kelly_fraction: float = KELLY_FRACTION, ): self.engine = FeatureEngine() self.ensemble = ModelEnsemble() self.calibrator = ProbabilityCalibrator() self.kelly = KellyCriterion(bankroll_usdc=bankroll_usdc, fraction=kelly_fraction) self.portfolio_kelly = PortfolioKelly(fraction=kelly_fraction) self.openmeteo = OpenMeteoClient() self.hko = HKOClient() self.spatial = SpatialWeatherClient() self.typhoon = TyphoonModel() self.min_edge_bps = min_edge_bps # Ensemble disagreement tracking self._last_forecast: Optional[pd.DataFrame] = None self._last_hourly: Optional[pd.DataFrame] = None self._last_features: Optional[np.ndarray] = None self._last_predictions: Optional[Dict[str, float]] = None self._last_spatial: Optional[Dict[str, float]] = None self._last_typhoon: Optional[Dict[str, float]] = None self._ensemble_spread: Optional[Dict[str, float]] = None # Load trained models self.models_loaded = self._load_models() def _load_models(self) -> bool: """Load trained models if available.""" try: self.ensemble.load_all() return len(self.ensemble.models) > 0 except Exception as e: print(f"ML models not loaded (train first): {e}") return False def fetch_and_predict(self, target_date: Optional[str] = None) -> Dict[str, float]: """ Fetch latest forecast and predict all targets. Returns dict of {target_name: calibrated_probability_0_100} """ # Fetch data daily = self.openmeteo.get_forecast(lead_days=7) if daily is None: print("MLPredictor: No forecast data available") return self._fallback_predictions() # Get hourly data from the client's internal cache hourly = getattr(self.openmeteo, '_last_hourly', None) self._last_forecast = daily self._last_hourly = hourly # Feature engineering X = self.engine.transform(daily, hourly) self._last_features = X # If models loaded, use ML predictions if self.models_loaded and len(self.ensemble.models) > 0: predictions = {} for day_idx in range(min(len(daily), 7)): date_str = daily.index[day_idx].strftime("%Y-%m-%d") if hasattr(daily.index[day_idx], 'strftime') else str(daily.index[day_idx]) if target_date and date_str != target_date and day_idx > 1: continue # Predict all loaded targets for this day X_day = X[day_idx].reshape(1, -1) day_probs = self.ensemble.predict_all(X_day) # Calibrate calibrated = {} for target, raw_prob in day_probs.items(): calibrated[target] = self.calibrator.calibrate(target, raw_prob) # Compute ensemble disagreement (multi-level) spread = self._compute_multimodel_spread(daily, hourly, day_idx) self._ensemble_spread = spread # Amplify edge based on spread for target in calibrated: adjusted = self._adjust_with_spread( calibrated[target], target, spread ) calibrated[target] = adjusted # Store date-str tagged predictions if day_idx <= 2: # Keep near-term predictions for target, prob in calibrated.items(): predictions[f"{target}_{date_str}"] = prob # Also store as raw target key (overwrites with latest) if day_idx == 1: # Tomorrow for target, prob in calibrated.items(): predictions[target] = prob self._last_predictions = predictions # Add typhoon predictions (not from LightGBM — separate model) self._last_typhoon = self._predict_typhoon() predictions.update(self._last_typhoon) # Add spatial features self._last_spatial = self._compute_spatial_features() return predictions # Fallback: use heuristic predictions return self._fallback_predictions() def _predict_typhoon(self) -> Dict[str, float]: """Generate typhoon signal-level probabilities.""" hko = self.hko current_signal = hko.get_current_signal_level() typhoon_info = hko.get_typhoon_info() month = datetime.now().month probs = {} for signal_level in ["T1", "T3", "T8"]: for lead_hours in [24, 48, 72, 120]: target = f"typhoon_{signal_level}_{lead_hours}h" prob = self.typhoon.signal_probability( signal_level=signal_level, lead_hours=lead_hours, current_signal=current_signal, current_conditions=typhoon_info, month=month, ) if lead_hours == 24: # Store short key too probs[f"typhoon_{signal_level}"] = prob probs[target] = prob return probs def _compute_spatial_features(self) -> Dict[str, float]: """Extract spatial features from multi-station forecasts.""" station_data = self.spatial.fetch_all_stations(lead_days=5) if not station_data: return {} return self.spatial.extract_spatial_features(station_data, day_index=1) def _fallback_predictions(self) -> Dict[str, float]: """Fallback heuristic predictions when no ML models loaded.""" if self._last_forecast is None or len(self._last_forecast) == 0: return {} d1 = self._last_forecast.iloc[min(1, len(self._last_forecast) - 1)] predictions = {} for target, tdef in TARGET_DEFINITIONS.items(): var = tdef["variable"] threshold = tdef["threshold"] if var in self._last_forecast.columns: val = float(d1.get(var, 0)) prob = self._heuristic_prob(val, threshold, target) predictions[target] = prob return predictions def _heuristic_prob(self, value: float, threshold: float, target: str) -> float: """Fallback heuristic: sigmoid-based probability.""" if "temp" in target: # Temperature: wider sigmoid, calibrated to HK summer excess = value - threshold return float(np.clip(50 + excess * 15, 3, 97)) elif "rain" in target: # Rain probabilities from Open-Meteo directly if threshold == 0: return float(np.clip(value, 0.5, 99.5)) else: return float(np.clip(value * 0.8 if threshold < 10 else value * 0.5, 1, 95)) elif "wind" in target: excess = value - threshold return float(np.clip(50 + excess * 5, 3, 97)) return 50.0 def _compute_multimodel_spread( self, daily: pd.DataFrame, hourly: pd.DataFrame, day_idx: int ) -> Dict[str, float]: """Compute ensemble disagreement metrics across model outputs. When multiple model outputs are available (GFS, ECMWF, WeatherNext), disagreement signifies uncertainty that the market may misprice. """ spread = {} # 1. Inter-day variability (persistence disagreement) if day_idx > 0 and len(daily) > day_idx: d0 = daily.iloc[day_idx - 1] d1 = daily.iloc[day_idx] spread["t2m_max_day_change"] = abs( float(d1.get("temperature_2m_max", 0)) - float(d0.get("temperature_2m_max", 0)) ) spread["precip_prob_day_change"] = abs( float(d1.get("precipitation_probability_max", 0)) - float(d0.get("precipitation_probability_max", 0)) ) spread["pressure_day_change"] = abs( float(hourly["surface_pressure"].mean() if "surface_pressure" in hourly else 1013) - float(hourly["surface_pressure"].iloc[max(0, day_idx * 24 - 24)] if "surface_pressure" in hourly else 1013) ) if hourly is not None and len(hourly) > 0 else 0.0 # 2. Wind direction variability (storm potential indicator) if hourly is not None and len(hourly) > 0 and "wind_direction_10m" in hourly: day_hourly = hourly.iloc[day_idx * 24:(day_idx + 1) * 24] if len(hourly) > (day_idx + 1) * 24 else hourly if len(day_hourly) > 0: wd = day_hourly["wind_direction_10m"].values spread["wind_dir_variance"] = float(np.var(wd)) if len(wd) > 1 else 0.0 # 3. Cloud structure complexity (convection proxy) if hourly is not None and len(hourly) > 0: for level in ["cloud_cover_low", "cloud_cover_mid", "cloud_cover_high"]: if level in hourly.columns: day_hourly = hourly.iloc[day_idx * 24:(day_idx + 1) * 24] if len(hourly) > (day_idx + 1) * 24 else hourly if len(day_hourly) > 0: spread[f"{level}_std"] = float(day_hourly[level].std()) # 4. Compute composite spread score (0-1) indicators = [] for k, v in spread.items(): if "t2m" in k: indicators.append(np.clip(v / 5.0, 0, 1)) # 5°C change = full signal elif "precip" in k: indicators.append(np.clip(v / 50.0, 0, 1)) # 50% change = full signal elif "pressure" in k: indicators.append(np.clip(v / 10.0, 0, 1)) # 10 hPa = full signal elif "variance" in k: indicators.append(np.clip(v / 5000.0, 0, 1)) elif "_std" in k: indicators.append(np.clip(v / 30.0, 0, 1)) spread["composite_spread"] = float(np.mean(indicators)) if indicators else 0.0 return spread def _adjust_with_spread( self, probability: float, target: str, spread: Dict[str, float] ) -> float: """ Adjust probability based on ensemble disagreement. When models disagree → higher uncertainty → wider confidence interval. In prediction markets, this often means the market price is LESS accurate (traders anchor on the wrong model or over-weight consensus). We amplify our edge when spread is high: push our probability further from 50% to reflect our confidence in the direction. """ composite = spread.get("composite_spread", 0.0) if composite < 0.1: return probability # Direction: is our prediction above or below 50%? direction = 1 if probability > 50 else -1 # Amplification: move probability up to spread * 20 bps further from 50 # High spread = more uncertainty = wider market spread = more edge amplification = min(composite * 20, 20) # Cap at 20 percentage points adjusted = probability + direction * amplification return float(np.clip(adjusted, 0.5, 99.5)) def generate_signal( self, target: str, market_probability: float, outcome: str = "YES", ) -> Dict: """ Generate a trading signal for a specific market. Parameters ---------- target : str Target name (e.g., 'temp_gt_30c_24h') market_probability : float Market-implied probability of the outcome (0-100) outcome : str Which outcome to bet on ('YES' or 'NO') Returns ------- Dict with model_prob, market_prob, edge_bps, kelly_size, side """ if not self._last_predictions: self.fetch_and_predict() model_prob = (self._last_predictions or {}).get(target, 50.0) cal_prob = self.calibrator.calibrate(target, model_prob) edge_bps = (cal_prob - market_probability) if abs(edge_bps) < self.min_edge_bps: return { "signal": "pass", "model_prob": cal_prob, "market_prob": market_probability, "edge_bps": edge_bps, "size_usdc": 0.0, } side = "buy_yes" if edge_bps > 0 else "buy_no" kelly_result = self.kelly.size_bet( our_probability=cal_prob, market_probability=market_probability, side=side, ) return { "signal": side, "model_prob": cal_prob, "market_prob": market_probability, "edge_bps": edge_bps, "size_usdc": kelly_result.size_usdc if kelly_result.kelly_active else 0.0, "kelly_fraction": kelly_result.fractional_kelly, "ensemble_spread": (self._ensemble_spread or {}).get("composite_spread", 0.0), } def record_outcome(self, target: str, predicted_prob: float, actual: bool): """Record resolved market outcome for calibration and correlation.""" self.calibrator.record_outcome( date=datetime.now().strftime("%Y-%m-%d"), variable=target, predicted_probability=predicted_prob, actual_outcome=actual, ) self.portfolio_kelly.record_outcomes({target: actual}) def get_top_signals( self, markets: List[Dict], default_market_prob: float = 50.0 ) -> List[Dict]: """Scan a list of market definitions and generate ranked signals.""" predictions = self._last_predictions or self.fetch_and_predict() signals = [] for market in markets: target = market.get("target", "") if target not in predictions: continue market_prob = market.get("market_probability", default_market_prob) signal = self.generate_signal(target, market_prob) if signal["signal"] != "pass": signals.append({ **signal, "target": target, "description": TARGET_DEFINITIONS.get(target, {}).get("description", ""), "question": market.get("question", ""), "condition_id": market.get("condition_id", ""), }) signals.sort(key=lambda s: abs(s["edge_bps"]), reverse=True) return signals def summary(self) -> str: """Human-readable summary of current predictions.""" predictions = self._last_predictions or {} lines = [] lines.append(f"\n=== ML Weather Predictions ({datetime.now():%Y-%m-%d %H:%M}) ===") lines.append(f" Models loaded: {len(self.ensemble.models)}") lines.append(f" Calibration records: {self.calibrator.get_calibration_stats(list(predictions.keys())[0] if predictions else 'temp_gt_30c_24h').get('n_observations', 0)}") if not predictions: lines.append(" No predictions available.") return "\n".join(lines) lines.append(f"\n Tomorrow's targets:") for target in TARGET_DEFINITIONS: if target in predictions: tdef = TARGET_DEFINITIONS[target] prob = predictions[target] lines.append(f" {tdef['description']}: {prob:.1f}%") if self._last_typhoon: lines.append(f"\n Typhoon probabilities:") for level in ["T1", "T3", "T8"]: key = f"typhoon_{level}" if key in self._last_typhoon: lines.append(f" {level}: {self._last_typhoon[key]:.1f}%") for h in [48, 72, 120]: for level in ["T1", "T8"]: key = f"typhoon_{level}_{h}h" if key in self._last_typhoon: lines.append(f" {level} in {h}h: {self._last_typhoon[key]:.1f}%") if self._last_spatial: uhi = self._last_spatial.get("uhi_tmax_delta", 0) instability = self._last_spatial.get("spatial_instability", 0) if abs(uhi) > 0.5 or instability > 0.2: lines.append(f"\n Spatial features:") if abs(uhi) > 0.5: lines.append(f" UHI delta: {uhi:+.1f}°C") if instability > 0.2: lines.append(f" Instability: {instability:.2f}") spread = self._ensemble_spread or {} if spread.get("composite_spread", 0) > 0.1: lines.append(f"\n Ensemble disagreement: {spread['composite_spread']:.2f} (amplified edge)") if spread.get("t2m_max_day_change", 0) > 0: lines.append(f" ΔTmax: {spread.get('t2m_max_day_change', 0):.1f}°C") return "\n".join(lines)