"""NautilusTrader strategy for HK weather prediction market trading. Integrates Open-Meteo + HKO weather forecasts with Polymarket CLOB execution. Flow: 1. On start: discover weather markets on Polymarket via Gamma API 2. Subscribe to L2 order book data for each discovered market 3. Every forecast_interval: run weather model, generate signal 4. On signal: compare model probability vs best bid/ask 5. Place limit order at favorable price when edge > threshold 6. On fill: track position, wait for resolution (0 or 1 payout) Instrument ID format: {condition_id}-{token_id}.POLYMARKET """ import asyncio import json from datetime import datetime, timedelta from typing import Optional, Dict, List import numpy as np import requests from nautilus_trader.cache.cache import Cache from nautilus_trader.common.component import Clock, LiveClock, MessageBus from nautilus_trader.config import StrategyConfig, ImportableStrategyConfig from nautilus_trader.core.uuid import UUID4 from nautilus_trader.live.node import TradingNode from nautilus_trader.model.book import OrderBook from nautilus_trader.model.data import QuoteTick, TradeTick, OrderBookDeltas, OrderBookDepth10 from nautilus_trader.model.enums import ( OrderSide, OrderType, TimeInForce, PositionSide, TriggerType, ) from nautilus_trader.model.events import OrderFilled from nautilus_trader.model.identifiers import ( ClientId, InstrumentId, PositionId, StrategyId, TraderId, VenueOrderId, ) from nautilus_trader.model.instruments import BinaryOption from nautilus_trader.model.objects import Price, Quantity, Money from nautilus_trader.model.orders import Order, LimitOrder, MarketOrder from nautilus_trader.model.position import Position from nautilus_trader.trading.strategy import Strategy from nautilus_trader.adapters.polymarket.common.constants import ( POLYMARKET_VENUE, POLYMARKET_CLIENT_ID, POLYMARKET_MAX_PRICE, POLYMARKET_MIN_PRICE, ) # FIXME: import from project package import sys sys.path.insert(0, "/home/satoshi/hk-weather-mkt") from weather.openmeteo_client import OpenMeteoClient from weather.hko_client import HKOClient from strategy.kelly import KellyCriterion GAMMA_API = "https://gamma-api.polymarket.com" class PolymarketWeatherStrategyConfig(StrategyConfig): """Configuration for the weather prediction market strategy.""" engine_type: type = Strategy bankroll_pusd: float = 1000.0 min_edge_bps: int = 200 max_position_per_market_pusd: float = 500.0 kelly_fraction: float = 0.25 forecast_interval_mins: int = 360 search_tags: tuple = ( "weather", "temperature", "hong kong", "typhoon", "precipitation", "climate", "heat", "storm", "rain", ) min_liquidity_usdc: float = 100.0 class PolymarketWeatherStrategy(Strategy): """Strategy that trades Polymarket weather outcome tokens using WeatherNext forecasts.""" def __init__(self, config: PolymarketWeatherStrategyConfig): super().__init__(config) self.config = config self.bankroll = config.bankroll_pusd self.kelly = KellyCriterion(bankroll_usdc=config.bankroll_pusd, fraction=config.kelly_fraction) # State self._instruments: Dict[InstrumentId, BinaryOption] = {} self._markets: Dict[str, Dict] = {} # condition_id -> market info self._best_bid: Dict[InstrumentId, float] = {} # instrument_id -> best bid self._best_ask: Dict[InstrumentId, float] = {} # instrument_id -> best ask self._positions: Dict[InstrumentId, Position] = {} self._orders: Dict[str, Order] = {} # client_order_id -> order self._active_signals: Dict[str, float] = {} # condition_id -> model probability # Weather clients self._openmeteo: Optional[OpenMeteoClient] = None self._hko: Optional[HKOClient] = None self._ml_predictor = None # ML predictor (lazy-loaded) self._last_forecast: Optional[Dict] = None self._last_ml_probs: Optional[Dict] = None # Task handles self._forecast_task: Optional[asyncio.Task] = None # ------------------------------------------------------------------- # # Lifecycle # # ------------------------------------------------------------------- # async def on_start(self): """Called when the strategy starts.""" self.log.info("Starting PolymarketWeatherStrategy") self._openmeteo = OpenMeteoClient() self._hko = HKOClient() try: from ml import MLPredictor self._ml_predictor = MLPredictor( bankroll_usdc=self.config.bankroll_pusd, min_edge_bps=self.config.min_edge_bps, kelly_fraction=self.config.kelly_fraction, ) self.log.info(f"ML predictor loaded: {len(self._ml_predictor.ensemble.models)} models") except Exception as e: self.log.warning(f"ML predictor not available: {e}") # Discover weather markets await self._discover_markets() if not self._instruments: self.log.warning("No weather markets found. Strategy will poll periodically.") # Start periodic forecast timer self._forecast_task = self.clock.loop.create_task(self._forecast_loop()) self.log.info(f"Forecast loop started (every {self.config.forecast_interval_mins}m)") async def on_stop(self): """Called when the strategy stops.""" if self._forecast_task: self._forecast_task.cancel() try: await self._forecast_task except asyncio.CancelledError: pass await self.cancel_all_orders(self.POLYMARKET_VENUE) self.log.info("Strategy stopped") async def on_instrument(self, instrument: BinaryOption): """Called when an instrument is loaded into the cache.""" self._instruments[instrument.id] = instrument self.log.info( f"Instrument: {instrument.id} " f"tick={instrument.price_increment} " f"min_qty={instrument.min_quantity} " f"max_qty={instrument.max_quantity}" ) async def on_disconnect(self): """Called when connection drops.""" self.log.warning("Disconnected from Polymarket. Reconnection in progress...") # ------------------------------------------------------------------- # # Market Data # # ------------------------------------------------------------------- # async def on_order_book_delta(self, deltas: OrderBookDeltas): """Order book updates.""" book = self.cache.order_book(deltas.instrument_id) if book: self._update_best_prices(book) async def on_order_book_depth10(self, depth: OrderBookDepth10): """Depth-10 snapshot.""" pass # Best prices already captured via deltas async def on_quote_tick(self, tick: QuoteTick): """Quote updates.""" pass async def on_trade_tick(self, tick: TradeTick): """Trade execution updates.""" pass # ------------------------------------------------------------------- # # Order Events # # ------------------------------------------------------------------- # async def on_order_filled(self, event: OrderFilled): """Order fill notification.""" order = event.to_order() self.log.info( f"FILLED {order.side} {order.quantity} @ {event.last_px} " f"[{order.instrument_id}] fee={event.commission}" ) # ------------------------------------------------------------------- # # Signal Generation # # ------------------------------------------------------------------- # async def _discover_markets(self): """Search Polymarket Gamma API for weather-related markets.""" self.log.info("Discovering weather markets on Polymarket...") discovered: Dict[str, Dict] = {} for tag in self.config.search_tags: try: params = { "tag": tag, "active": "true", "closed": "false", "limit": 50, "order": "liquidity", } resp = requests.get(f"{GAMMA_API}/markets", params=params, timeout=15) resp.raise_for_status() for m in resp.json(): cid = m.get("conditionId") if not cid or cid in discovered: continue liquidity = float(m.get("liquidity", 0)) if liquidity < self.config.min_liquidity_usdc: continue discovered[cid] = { "condition_id": cid, "question": m.get("question", ""), "slug": m.get("slug", ""), "volume": float(m.get("volume", 0)), "liquidity": liquidity, "end_date": m.get("endDateIso", ""), "tag": tag, } except Exception as e: self.log.warning(f"Gamma API error for tag '{tag}': {e}") self._markets = discovered self.log.info(f"Found {len(discovered)} weather-related markets") # List top markets sorted_mkts = sorted(discovered.values(), key=lambda m: m["liquidity"], reverse=True) for m in sorted_mkts[:10]: self.log.info( f" [{m['liquidity']:.0f} USDC liq] {m['question'][:80]} " f"(tag={m['tag']})" ) # Subscribe to order books for discovered markets # We need to load instruments first, then subscribe for m in sorted_mkts: try: instruments = await self._load_instruments_for_condition(m["condition_id"]) for inst in instruments: self.subscribe_order_book_deltas(inst.id) self.log.info(f" Subscribed to {inst.id}") except Exception as e: self.log.warning(f" Failed to load instruments for {m['condition_id']}: {e}") async def _load_instruments_for_condition(self, condition_id: str) -> List: """Load BinaryOption instruments for a Polymarket condition.""" # We need to query the CLOB API for the market's tokens # The instrument provider handles this instruments = [] try: clob_resp = requests.get( f"https://clob.polymarket.com/markets/{condition_id}", timeout=10, ) clob_resp.raise_for_status() clob_data = clob_resp.json() tokens = clob_data.get("tokens", []) for token in tokens: token_id = token.get("token_id") if token_id: instrument = self.cache.instrument( InstrumentId.from_str(f"{condition_id}-{token_id}.POLYMARKET") ) if instrument: instruments.append(instrument) except Exception as e: self.log.warning(f"Failed to load instruments for {condition_id}: {e}") return instruments async def _forecast_loop(self): """Periodically run weather forecast and generate trading signals.""" while True: try: await self._update_forecast() await self._generate_signals() except Exception as e: self.log.error(f"Forecast loop error: {e}") await asyncio.sleep(self.config.forecast_interval_mins * 60) async def _update_forecast(self): """Fetch latest weather forecast data.""" self.log.info("Updating weather forecast...") try: # Use ML predictor if available if self._ml_predictor and self._ml_predictor.models_loaded: self._ml_predictor.fetch_and_predict() self._last_ml_probs = self._ml_predictor._last_predictions self.log.info(f"ML forecast: {self._ml_predictor.summary()}") tomorrow = (datetime.now() + timedelta(days=1)).strftime("%Y-%m-%d") self._last_forecast = self._openmeteo.get_scoring_window_summary(tomorrow) hko_fc = self._hko.get_forecast() if hko_fc: self._last_forecast["hko_tomorrow"] = hko_fc[0] if hko_fc else None typhoon = self._hko.get_typhoon_info() if typhoon: self._last_forecast["typhoon"] = typhoon current = self._hko.get_current_weather() if current: temps = current.get("temperature", []) self._last_forecast["current_temp"] = temps[0]["value"] if temps else None except Exception as e: self.log.error(f"Forecast fetch error: {e}") async def _generate_signals(self): """Generate trading signals by comparing model forecast vs market prices.""" if not self._last_forecast: self.log.warning("No forecast data available for signal generation") return for condition_id, market in self._markets.items(): question = market["question"].lower() model_prob = self._compute_model_probability(question) if model_prob is None: continue self._active_signals[condition_id] = model_prob # Get current market price (mid of best bid/ask) yes_instrument_id = InstrumentId.from_str(f"{condition_id}-yes.POLYMARKET") no_instrument_id = InstrumentId.from_str(f"{condition_id}-no.POLYMARKET") # Determine which side to bet # The YES token is the one we buy if we think the event WILL happen mid_price = self._get_mid_price(yes_instrument_id) if mid_price is None: # Try NO token price as alternative mid_price_no = self._get_mid_price(no_instrument_id) if mid_price_no is not None: mid_price = 1.0 - mid_price_no # P(YES) = 1 - P(NO) if mid_price is None: continue market_prob = mid_price * 100.0 # Convert to percentage edge_bps = (model_prob - market_prob) * 100.0 if abs(edge_bps) < self.config.min_edge_bps: continue # Kelly sizing side = "buy_yes" if edge_bps > 0 else "buy_no" kelly_result = self.kelly.size_bet( our_probability=model_prob, market_probability=market_prob, side=side, max_size=self.config.max_position_per_market_pusd, ) if not kelly_result.kelly_active or kelly_result.size_usdc < 1.0: continue await self._place_weather_order( condition_id=condition_id, question=market["question"], side=side, kelly=kelly_result, model_prob=model_prob, market_prob=market_prob, edge_bps=edge_bps, ) def _compute_model_probability(self, question: str) -> Optional[float]: """Compute our model's probability for a given market question. Uses ML model predictions when available, falls back to heuristics. Also records resolved outcomes for calibration. """ # Try ML predictor first if self._ml_predictor and self._last_ml_probs: target = self._question_to_target(question) if target and target in self._last_ml_probs: return self._last_ml_probs[target] # Fallback heuristic if not self._last_forecast: return None return self._compute_heuristic_probability(question) @staticmethod def _question_to_target(question: str) -> Optional[str]: """Map a Polymarket question to an ML model target.""" q = question.lower() if "rain" in q or "precipitation" in q: if "10mm" in q or "10 mm" in q or "heavy" in q: return "rain_gt_10mm_24h" if "5mm" in q or "5 mm" in q: return "rain_gt_5mm_24h" return "rain_gt_0mm_24h" if "temperature" in q or "temp" in q: if "35" in q or "thirty five" in q: return "temp_gt_35c_24h" if "33" in q or "thirty three" in q: return "temp_gt_33c_24h" if "30" in q or "thirty" in q: return "temp_gt_30c_24h" return "temp_gt_30c_24h" if "wind" in q or "gust" in q: return "wind_gt_30kmh_24h" if "typhoon" in q or "t8" in q or "cyclone" in q: return None # No ML model for typhoon yet return None def _compute_heuristic_probability(self, question: str) -> Optional[float]: """Fallback heuristic probability (legacy).""" if not self._last_forecast: return None question = question.lower() if "rain" in question or "precipitation" in question: return self._last_forecast.get("precipitation_probability_max", None) if "temperature" in question and "above" in question: tmax = self._last_forecast.get("temperature_2m_max", 30) if "30" in question or "thirty" in question: threshold = 30.0 elif "35" in question or "thirty five" in question: threshold = 35.0 elif "33" in question or "thirty three" in question: threshold = 33.0 elif "40" in question or "forty" in question: threshold = 40.0 else: return None return min(97.0, max(3.0, 50.0 + (tmax - threshold) * 20.0)) if "typhoon" in question or "t8" in question or "tropical cyclone" in question: typhoon = self._last_forecast.get("typhoon", None) return 30.0 if typhoon else 5.0 if "heat" in question or "hot" in question: tmax = self._last_forecast.get("temperature_2m_max", 30) return min(97.0, max(3.0, 50.0 + (tmax - 33.0) * 25.0)) # Default: use rain probability for general weather questions return self._last_forecast.get("precipitation_probability_max", 50.0) def _get_mid_price(self, instrument_id: InstrumentId) -> Optional[float]: """Get mid-price from order book.""" try: book = self.cache.order_book(instrument_id) if not book: return None if book.best_bid_price() and book.best_ask_price(): bid = book.best_bid_price().as_f64() ask = book.best_ask_price().as_f64() return (bid + ask) / 2.0 elif book.best_bid_price(): return book.best_bid_price().as_f64() elif book.best_ask_price(): return book.best_ask_price().as_f64() except Exception: pass return None # ------------------------------------------------------------------- # # Order Placement # # ------------------------------------------------------------------- # async def _place_weather_order( self, condition_id: str, question: str, side: str, kelly: "KellyResult", model_prob: float, market_prob: float, edge_bps: float, ): """Place a limit order on Polymarket based on weather signal.""" token_id = "yes" if side == "buy_yes" else "no" instrument_id = InstrumentId.from_str(f"{condition_id}-{token_id}.POLYMARKET") instrument = self.cache.instrument(instrument_id) if not instrument: self.log.warning(f"Instrument not in cache: {instrument_id}") return # Our limit price = model-implied fair value # If we think YES prob is 65% and market at 50%, we bid 0.55 (midway) our_price = model_prob / 100.0 if side == "buy_yes" else (100.0 - model_prob) / 100.0 # Use Kelly edge to set aggressive but fair price # Buy at a price between market and our fair value market_price = market_prob / 100.0 if side == "buy_yes" else (100.0 - market_prob) / 100.0 limit_price = (our_price + market_price) / 2.0 # Clamp to venue bounds tick_size = instrument.price_increment limit_price = max( POLYMARKET_MIN_PRICE, min(POLYMARKET_MAX_PRICE, limit_price), ) # Round to tick size limit_price = round(limit_price / tick_size) * tick_size limit_price = max(tick_size, min(1.0 - tick_size, limit_price)) # Convert pUSD notional to share quantity # shares = pUSD / price (for YES), pUSD / (1-price) (for NO) if side == "buy_yes": shares = kelly.size_usdc / max(limit_price, 0.0001) else: shares = kelly.size_usdc / max(1.0 - limit_price, 0.0001) # Round shares to 2 decimal places (Polymarket precision) shares = round(shares, 2) if shares < 0.01: self.log.info(f"Order too small: {shares} shares") return # Build limit order price = Price(limit_price, instrument.price_precision) qty = Quantity(shares, instrument.size_precision) order = self.order_factory.limit( instrument_id=instrument_id, order_side=OrderSide.BUY if side == "buy_yes" else OrderSide.SELL, quantity=qty, price=price, time_in_force=TimeInForce.GTC, post_only=True, # Maker orders: no taker fees ) self.submit_order(order, position_id=None) self.log.info( f"ORDER: {side.upper()} {shares} shares @ {limit_price:.4f} " f"'{question[:60]}' " f"(model={model_prob:.1f}%, mkt={market_prob:.1f}%, " f"edge={'+' if edge_bps > 0 else ''}{edge_bps:.0f}bps)" ) # ------------------------------------------------------------------- # # Helpers # # ------------------------------------------------------------------- # def _update_best_prices(self, book: OrderBook): """Track best bid/ask from order book updates.""" if book.best_bid_price(): self._best_bid[book.instrument_id] = book.best_bid_price().as_f64() if book.best_ask_price(): self._best_ask[book.instrument_id] = book.best_ask_price().as_f64()