Files
ramseshk 7d7a67bd20 Add ML prediction pipeline — LightGBM, calibration fix, ensemble disagreement
Tier 1 ML enhancements:
- Feature engineering (37 features across 5 groups: thermal, dynamic,
  moisture, temporal, interaction) from NWP model output
- 7 LightGBM probability models for rain/temp/wind thresholds
- Temperature-scaled probabilities to prevent overconfidence on bootstrap data
- MLPredictor: unified inference pipeline replacing heuristic sigmoids
- Ensemble disagreement signals (composite spread → edge amplification)
- Fixed calibration loop: update_calibration() now functional (EMA of errors)
- record_outcome() wired for post-resolution feedback
- Nautilus strategy updated: ML predictions take priority, heuristics as fallback
- Historical backtest engine with Sharpe/ROI/max-DD simulation
- Bootstrap training data generator from HK climate normals

Run: python ml/train.py && python ml/backtest.py --edge 50
2026-08-10 17:50:07 +08:00

577 lines
23 KiB
Python

"""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()