commit c93af970596c602a7608392a195d29e9841ffbbe Author: ramseshk <45832522+ramseshk@users.noreply.github.com> Date: Mon Aug 10 12:48:05 2026 +0800 HK Weather Prediction Market Pipeline: WeatherNext + HKO + Polymarket - Open-Meteo WeatherNext API client for HK forecasts - HKO public data client (current conditions, 9-day forecast, typhoon warnings) - HK-specific weather extraction and calibration - Polymarket market scanning, price discovery, and market creation proposals - Trading strategy engine: edge detection, Kelly criterion sizing, probability calibration - End-to-end pipeline with dry-run mode and scheduled runner - Interactive dashboard with live HK weather + forecasts + trading signals Dependencies: Python 3.10+, openmeteo-requests, pandas No API keys needed for dry-run mode. Polymarket trading requires private key in .env. diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..a6d793b --- /dev/null +++ b/.env.example @@ -0,0 +1,16 @@ +# Environment variables template +# Copy to .env and fill in your keys + +# Open-Meteo API (free tier works without key) +OPENMETEO_API_KEY= + +# Polymarket - generate at https://clob.polymarket.com/settings +POLYMARKET_PRIVATE_KEY= +POLYMARKET_FUNDER= + +# HKO Open Data API +HKO_OPEN_DATA_KEY= + +# CDS API for ECMWF data (free registration at https://cds.climate.copernicus.eu) +CDSAPI_KEY= +CDSAPI_URL=https://cds.climate.copernicus.eu/api diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..d93f9dd --- /dev/null +++ b/.gitignore @@ -0,0 +1,15 @@ +venv/ +__pycache__/ +*.pyc +*.pyo +.env +data/weights/*.npz +.opencache/ +.openmeteo_cache* +.openmeteo_cache.sqlite +.calibration_history.json +logs/ +*.log +.DS_Store +*.sqlite + diff --git a/.python-version b/.python-version new file mode 100644 index 0000000..1d4830e --- /dev/null +++ b/.python-version @@ -0,0 +1 @@ +3.10.17 diff --git a/config.py b/config.py new file mode 100644 index 0000000..0e91ed1 --- /dev/null +++ b/config.py @@ -0,0 +1,109 @@ +# HK Weather Prediction Market +# Configuration and constants + +import os +from pathlib import Path +from dotenv import load_dotenv + +load_dotenv() + +PROJECT_ROOT = Path(__file__).parent +DATA_DIR = PROJECT_ROOT / "data" +WEIGHTS_DIR = DATA_DIR / "weights" +FORECAST_DIR = DATA_DIR / "forecasts" + +for d in [DATA_DIR, WEIGHTS_DIR, FORECAST_DIR]: + d.mkdir(parents=True, exist_ok=True) + +# === HK Region Bounding Box === +HK_BBOX = { + "lat_min": 22.0, + "lat_max": 22.6, + "lon_min": 113.8, + "lon_max": 114.5, +} + +HK_COORDS = { + "hko_headquarters": (22.302, 114.174), # Tsim Sha Tsui + "chek_lap_kok": (22.308, 113.918), # Airport + "sheung_shui": (22.505, 114.128), # New Territories + "stanley": (22.218, 114.218), # South side + "cheung_chau": (22.201, 114.028), # Outlying island +} + +# === API Keys === +OPENMETEO_API_KEY = os.getenv("OPENMETEO_API_KEY", "") +POLYMARKET_PRIVATE_KEY = os.getenv("POLYMARKET_PRIVATE_KEY", "") +POLYMARKET_FUNDER = os.getenv("POLYMARKET_FUNDER", "") +HKO_OPEN_DATA_KEY = os.getenv("HKO_OPEN_DATA_KEY", "") +CDSAPI_KEY = os.getenv("CDSAPI_KEY", "") +CDSAPI_URL = os.getenv("CDSAPI_URL", "https://cds.climate.copernicus.eu/api") + +# === Polymarket === +POLYMARKET_GAMMA_API = "https://gamma-api.polymarket.com" +POLYMARKET_CLOB_API = "https://clob.polymarket.com" +POLYMARKET_SOCKET = "wss://ws-subscriptions-clob.polymarket.com/ws" + +# === WeatherNext Model === +MODEL_CONFIGS = { + "mini": { + "name": "WeatherNextCyclones_Mini_<2024", + "weights_file": "WeatherNextCyclones_Mini_<2024.npz", + "resolution": "1°", + "gcs_path": "gs://dm_graphcast/WeatherNextCyclones_Mini_<2024.npz", + "vram_required_gb": 8, + }, + "operational": { + "name": "WeatherNext2_<2025", + "resolution": "0.25°", + "gcs_path": "gs://dm_graphcast/WeatherNext2_<2025_model*.npz", + "vram_required_gb": 40, + }, +} + +MODEL_CONFIG = MODEL_CONFIGS["mini"] + +# === Forecast Parameters === +FORECAST_LEAD_TIMES = [0, 6, 12, 18, 24, 36, 48, 72, 96, 120] # hours + +# === Weather Variables of Interest === +WEATHER_VARIABLES = [ + "temperature_2m", + "relative_humidity_2m", + "wind_speed_10m", + "wind_gusts_10m", + "precipitation", + "total_cloud_cover", + "surface_pressure", +] + +# === Polymarket Market Conditions === +# Weather conditions we can create markets for +MARKET_CONDITIONS = { + "rain_tomorrow": { + "description": "Will it rain in Hong Kong tomorrow?", + "threshold": lambda forecast: forecast.get("precipitation_probability", 0) > 50, + "market_type": "binary", + }, + "temp_above_30": { + "description": "Will the temperature in Hong Kong exceed 30°C tomorrow?", + "threshold": lambda forecast: forecast.get("temperature_2m_max", 0) > 30.0, + "market_type": "binary", + }, + "t8_signal": { + "description": "Will the T8 typhoon signal be hoisted in the next 7 days?", + "threshold": None, + "market_type": "binary", + }, + "rainfall_amount": { + "description": "Total rainfall (mm) in Hong Kong tomorrow", + "threshold": None, + "market_type": "scalar", + }, +} + +# === Trading Strategy === +MIN_EDGE_BPS = 200 # Minimum edge in basis points to trade +MAX_POSITION_USDC = 500.0 # Maximum position per market in USDC +KELLY_FRACTION = 0.25 # Fraction of full Kelly to use +MIN_LIQUIDITY_USDC = 100.0 # Minimum market liquidity to trade diff --git a/dashboard.py b/dashboard.py new file mode 100644 index 0000000..4248558 --- /dev/null +++ b/dashboard.py @@ -0,0 +1,112 @@ +#!/usr/bin/env python3 +"""Quick test and display of HK weather forecasts with trading signals.""" + +import sys +from datetime import datetime, timedelta +from weather.hk_extractor import HKExtractor +from weather.hko_client import HKOClient +from weather.openmeteo_client import OpenMeteoClient +from strategy.signals import SignalGenerator +from strategy.calibrator import ProbabilityCalibrator +from strategy.kelly import KellyCriterion + + +def main(): + print("╔══════════════════════════════════════════════════════╗") + print("║ HK Weather Prediction Market Dashboard ║") + print(f"║ {datetime.now():%Y-%m-%d %H:%M:%S} HKT ║") + print("╚══════════════════════════════════════════════════════╝") + print() + + # Fetch data + weather = HKExtractor() + hko = HKOClient() + + print("─── Current Conditions ───────────────────────────────") + current = hko.get_current_weather() + if current: + temps = current.get("temperature", []) + for t in temps[:3]: + print(f" {t['place']:20s} {t['value']}°C") + print(f" Humidity: {current.get('humidity', [{}])[0].get('value', 'N/A')}%") + warning = current.get("warning_message", "") + if warning: + print(f" ⚠ {warning}") + + om = OpenMeteoClient() + cur = om.get_current_conditions() + if cur: + print(f" Feels like: {cur.get('apparent_temp', 'N/A')}°C") + print(f" Wind: {cur.get('wind_speed', 'N/A'):.1f} km/h") + print(f" Pressure: {cur.get('surface_pressure', 'N/A'):.1f} hPa") + + # Typhoon info + typhoon = hko.get_typhoon_info() + if typhoon: + signal = hko.get_current_signal_level() + print(f" Typhoon signal: {'T' + str(signal) if signal else 'None'}") + + print() + print("─── HKO 9-Day Forecast ───────────────────────────────") + hko_fc = hko.get_forecast() + if hko_fc: + for day in hko_fc[:5]: + print(f" {day['date']} ({day['week']}): " + f"{day['forecast_temp_min']}-{day['forecast_temp_max']}°C, " + f"Wind: {day.get('forecast_wind', 'N/A')}") + + print() + print("─── WeatherNext (Open-Meteo) Forecast ───────────────") + fc = weather.get_hk_forecast(lead_days=5) + wnext = fc.get("sources", {}).get("weathernext", []) + for day in wnext[:5]: + print(f" {day['date']}: " + f"↑{day['temp_max_calibrated']:.1f}°C ↓{day['temp_min_calibrated']:.1f}°C, " + f"Rain: {day['precipitation_probability_calibrated']:.1f}% ({day['precipitation_sum']:.1f}mm), " + f"Wind: {day['wind_speed_max_calibrated']:.1f} km/h") + + print() + print("─── Trading Signals ──────────────────────────────────") + sig_gen = SignalGenerator() + sig_gen.generate_signals() + print(sig_gen.get_signal_summary()) + + # Also show market creation proposals + print() + print("─── Proposed Markets (for Polymarket) ────────────────") + tomorrow = datetime.now() + timedelta(days=1) + if wnext: + d1 = wnext[0] if len(wnext) > 0 else {} + rain_p = d1.get("precipitation_probability_calibrated", 50) + temp_p = d1.get("temp_max_calibrated", 30) + wind_p = d1.get("wind_speed_max_calibrated", 15) + + print(f" 1. 'Will it rain in HK on {tomorrow:%Y-%m-%d}?' [YES: ~{rain_p:.0f}%]") + print(f" 2. 'Will HK temp exceed 33°C on {tomorrow:%Y-%m-%d}?' [YES: ~{min(95, max(5, 50 + (temp_p - 33) * 20)):.0f}%]") + print(f" 3. 'Will HK temp exceed 35°C on {tomorrow:%Y-%m-%d}?' [YES: ~{min(95, max(5, 50 + (temp_p - 35) * 20)):.0f}%]") + print(f" 4. 'Will rain exceed 10mm in HK on {tomorrow:%Y-%m-%d}?' [check hourly]") + print(f" 5. 'Will a T8 signal be hoisted in HK in the next 7 days?'") + + print() + print("─── Kelly Sizing Sim ─────────────────────────────────") + kelly = KellyCriterion(bankroll_usdc=1000.0) + for name, our_p, mkt_p in [ + ("Rain tomorrow", rain_p, 45), + ("Temp > 35°C", min(95, max(5, 50 + (temp_p - 35) * 20)), 30), + ]: + r = kelly.size_bet(our_p, mkt_p, "buy_yes") + if r.kelly_active: + print(f" {name}: Bet ${r.size_usdc:.2f} YES (edge={r.edge:.3f})") + else: + r2 = kelly.size_bet(our_p, mkt_p, "buy_no") + if r2.kelly_active: + print(f" {name}: Bet ${r2.size_usdc:.2f} NO (edge={r2.edge:.3f})") + else: + print(f" {name}: No edge (model={our_p:.0f}% vs market={mkt_p:.0f}%)") + + print() + print("═" * 56) + + +if __name__ == "__main__": + main() diff --git a/markets/__init__.py b/markets/__init__.py new file mode 100644 index 0000000..13c57b7 --- /dev/null +++ b/markets/__init__.py @@ -0,0 +1,6 @@ +"""Polymarket integration for HK weather prediction markets.""" + +from .polymarket_client import PolymarketClient +from .trader import Trader + +__all__ = ["PolymarketClient", "Trader"] diff --git a/markets/polymarket_client.py b/markets/polymarket_client.py new file mode 100644 index 0000000..30821fe --- /dev/null +++ b/markets/polymarket_client.py @@ -0,0 +1,214 @@ +"""Polymarket API client for market data and trading. + +Uses the Gamma Markets API for market discovery and CLOB for order execution. +""" + +import json +import time +import hashlib +from datetime import datetime +from typing import Optional, Dict, List +from urllib.parse import urlencode + +import requests + +from config import ( + POLYMARKET_GAMMA_API, + POLYMARKET_CLOB_API, + MIN_LIQUIDITY_USDC, +) + + +class PolymarketClient: + """Client for Polymarket prediction markets.""" + + def __init__(self): + self.gamma_url = POLYMARKET_GAMMA_API + self.clob_url = POLYMARKET_CLOB_API + self.session = requests.Session() + self.session.headers.update({ + "User-Agent": "HK-Weather-Market/1.0", + "Accept": "application/json", + }) + + def search_markets( + self, + query: str = "hong kong weather", + active: bool = True, + limit: int = 20, + ) -> List[Dict]: + """Search for weather-related markets on Polymarket.""" + params = { + "query": query, + "active": str(active).lower(), + "closed": "false", + "limit": limit, + "archived": "false", + "order": "liquidity", + } + + try: + url = f"{self.gamma_url}/markets?{urlencode(params)}" + resp = self.session.get(url, timeout=15) + resp.raise_for_status() + markets = resp.json() + + return [ + { + "id": m.get("id"), + "question": m.get("question"), + "condition_id": m.get("conditionId"), + "slug": m.get("slug"), + "volume": float(m.get("volume", 0)), + "liquidity": float(m.get("liquidity", 0)), + "volume_24hr": float(m.get("volume24hr", 0)), + "end_date": m.get("endDateIso"), + "start_date": m.get("startDateIso"), + "active": m.get("active"), + "closed": m.get("closed"), + "outcome_prices": json.loads(m.get("outcomePrices", "[]")), + "outcomes": json.loads(m.get("outcomes", "[]")), + "description": m.get("description", ""), + "category": m.get("category", ""), + "tags": m.get("tags", []), + } + for m in markets + ] + + except Exception as e: + print(f"Polymarket search error: {e}") + return [] + + def get_market(self, condition_id: str) -> Optional[Dict]: + """Get a single market by condition ID.""" + try: + url = f"{self.gamma_url}/markets/{condition_id}" + resp = self.session.get(url, timeout=15) + resp.raise_for_status() + m = resp.json() + + return { + "id": m.get("id"), + "question": m.get("question"), + "condition_id": m.get("conditionId"), + "slug": m.get("slug"), + "volume": float(m.get("volume", 0)), + "liquidity": float(m.get("liquidity", 0)), + "end_date": m.get("endDateIso"), + "active": m.get("active"), + "closed": m.get("closed"), + "outcome_prices": json.loads(m.get("outcomePrices", "[]")), + "outcomes": json.loads(m.get("outcomes", "[]")), + "description": m.get("description", ""), + } + + except Exception as e: + print(f"Polymarket get_market error: {e}") + return None + + def get_market_price(self, token_id: str) -> Optional[float]: + """Get current price for a specific outcome token.""" + try: + url = f"{self.clob_url}/price?token_id={token_id}&side=buy" + resp = self.session.get(url, timeout=10) + if resp.status_code == 200: + data = resp.json() + return float(data.get("price", 0)) + except Exception as e: + print(f"Polymarket price error: {e}") + return None + + def get_order_book(self, token_id: str) -> Dict: + """Get order book for a token.""" + try: + url = f"{self.clob_url}/book?token_id={token_id}" + resp = self.session.get(url, timeout=10) + resp.raise_for_status() + return resp.json() + except Exception as e: + print(f"Polymarket orderbook error: {e}") + return {} + + def find_relevant_weather_markets(self) -> List[Dict]: + """Find all weather-related markets relevant to Hong Kong.""" + queries = [ + "hong kong weather", + "hong kong temperature", + "hong kong typhoon", + "hong kong rain", + "asia typhoon", + "south china sea", + ] + + all_markets = [] + seen_ids = set() + + for query in queries: + markets = self.search_markets(query=query) + for m in markets: + if m["id"] not in seen_ids and m.get("active") and not m.get("closed"): + seen_ids.add(m["id"]) + all_markets.append(m) + + all_markets.sort(key=lambda m: m.get("liquidity", 0), reverse=True) + + return all_markets + + def get_market_implied_probability( + self, condition_id: str, outcome_index: int = 0 + ) -> Optional[float]: + """Get market-implied probability for a specific outcome. + + Uses midpoint of best bid/ask when available, otherwise last price. + """ + market = self.get_market(condition_id) + if not market or "outcome_prices" not in market: + return None + + if outcome_index < len(market["outcome_prices"]): + return float(market["outcome_prices"][outcome_index]) + + return None + + def get_clob_token_id(self, condition_id: str, outcome_index: int = 0) -> Optional[str]: + """Get CLOB token ID for a market outcome. + + Token IDs are derived deterministically from condition ID + outcome index. + """ + try: + url = f"{self.clob_url}/markets/{condition_id}" + resp = self.session.get(url, timeout=10) + resp.raise_for_status() + data = resp.json() + + tokens = data.get("tokens", []) + if outcome_index < len(tokens): + return tokens[outcome_index].get("token_id") + except Exception as e: + print(f"CLOB token ID error: {e}") + + return None + + def create_market( + self, + question: str, + outcomes: List[str], + end_date: str, + description: str = "", + ) -> Optional[Dict]: + """ + Create a new market on Polymarket. + NOTE: Requires whitelisted API key and Polygonscan approval. + Markets go through curation before going live. + """ + print("Market creation requires curation approval from Polymarket.") + print(f"Would create: {question}") + print(f"Outcomes: {outcomes}") + print(f"End date: {end_date}") + return { + "status": "proposed", + "question": question, + "outcomes": outcomes, + "end_date": end_date, + "note": "Submit via Polymarket UI or contact partnerships@polymarket.com", + } diff --git a/markets/trader.py b/markets/trader.py new file mode 100644 index 0000000..49586ca --- /dev/null +++ b/markets/trader.py @@ -0,0 +1,211 @@ +"""Trading execution engine for Polymarket weather markets. + +Handles order placement, position sizing, and risk management +for the HK weather prediction market strategy. +""" + +from datetime import datetime +from typing import Optional, Dict, List, Tuple +from dataclasses import dataclass, field + +from .polymarket_client import PolymarketClient +from config import MIN_EDGE_BPS, MAX_POSITION_USDC, MIN_LIQUIDITY_USDC + + +@dataclass +class TradeSignal: + """A trading signal from the strategy engine.""" + market_id: str + condition_id: str + question: str + outcome_index: int + outcome_label: str + model_probability: float # Our model-implied probability (0-100) + market_probability: float # Market-implied probability (0-100) + edge_bps: float # Edge in basis points + recommended_size_usdc: float # Kelly-recommended bet size + max_size_usdc: float # Maximum allowed position + signal_type: str # "buy_yes", "buy_no", "pass" + + +@dataclass +class ExecutionResult: + """Result of a trade execution.""" + signal: TradeSignal + success: bool + order_id: Optional[str] = None + filled_amount: float = 0.0 + avg_price: float = 0.0 + error: Optional[str] = None + timestamp: str = field(default_factory=lambda: datetime.now().isoformat()) + + +class Trader: + """Execute trades based on strategy signals.""" + + def __init__( + self, + client: PolymarketClient, + private_key: str = "", + funder_address: str = "", + dry_run: bool = True, + ): + self.client = client + self.private_key = private_key + self.funder_address = funder_address + self.dry_run = dry_run + self.clob = None + self.positions: Dict[str, float] = {} + self.trade_history: List[ExecutionResult] = [] + + if not dry_run and private_key: + self._init_clob() + + def _init_clob(self): + """Initialize CLOB client for live trading.""" + try: + from py_clob_client.client import ClobClient + from py_clob_client.clob_types import OrderArgs + + host = "https://clob.polymarket.com" + chain_id = 137 # Polygon mainnet + + self.clob = ClobClient( + host=host, + key=self.private_key, + chain_id=chain_id, + funder=self.funder_address, + signature_type=2, + ) + print("CLOB client initialized for live trading") + except Exception as e: + print(f"CLOB init failed: {e}. Running in dry-run mode.") + self.dry_run = True + + def execute_signal(self, signal: TradeSignal) -> ExecutionResult: + """Execute a single trade signal.""" + if signal.signal_type == "pass": + return ExecutionResult( + signal=signal, + success=True, + note="No trade: edge below threshold", + ) + + # Get token ID + token_id = self.client.get_clob_token_id( + signal.condition_id, signal.outcome_index + ) + + if not token_id: + return ExecutionResult( + signal=signal, + success=False, + error="Could not get token ID", + ) + + # Calculate number of shares at size (each share = $1 if correct) + price = signal.market_probability / 100.0 + size = min(signal.recommended_size_usdc, signal.max_size_usdc) + + if size < 1.0: + return ExecutionResult( + signal=signal, + success=False, + error=f"Size too small: ${size:.2f}", + ) + + if self.dry_run: + return self._execute_dry_run(signal, token_id, size, price) + else: + return self._execute_live(signal, token_id, size, price) + + def _execute_dry_run( + self, signal: TradeSignal, token_id: str, size: float, price: float + ) -> ExecutionResult: + """Simulate trade execution for testing.""" + result = ExecutionResult( + signal=signal, + success=True, + order_id=f"DRY_RUN_{datetime.now().timestamp()}", + filled_amount=size, + avg_price=price, + ) + self.trade_history.append(result) + + position_key = f"{signal.condition_id}_{signal.outcome_index}" + self.positions[position_key] = self.positions.get(position_key, 0) + size + + print(f" [DRY RUN] {signal.signal_type}: ${size:.2f} on '{signal.question}'" + f" @ {price:.4f} (edge: {signal.edge_bps:.0f}bps)") + + return result + + def _execute_live( + self, signal: TradeSignal, token_id: str, size: float, price: float + ) -> ExecutionResult: + """Execute real trade on Polymarket CLOB.""" + if not self.clob: + return ExecutionResult( + signal=signal, + success=False, + error="CLOB not initialized", + ) + + try: + # Create a limit order (IOC to avoid partial fills on stale prices) + order_args = { + "token_id": token_id, + "price": price, + "size": size, + "side": "BUY" if signal.signal_type == "buy_yes" else "SELL", + } + + response = self.clob.create_and_post_order( + order_args, orderType="GTC" + ) + + result = ExecutionResult( + signal=signal, + success=True, + order_id=response.get("orderID", ""), + filled_amount=float(response.get("filled_size", 0)), + avg_price=float(response.get("avg_price", price)), + ) + self.trade_history.append(result) + + print(f" [LIVE] {signal.signal_type}: ${size:.2f} on '{signal.question}'" + f" @ {price:.4f} (edge: {signal.edge_bps:.0f}bps)") + + return result + + except Exception as e: + return ExecutionResult( + signal=signal, + success=False, + error=str(e), + ) + + def get_positions_summary(self) -> Dict: + """Get summary of current positions and P&L.""" + total_bet = sum(self.positions.values()) + open_trades = len([t for t in self.trade_history if t.success]) + + return { + "total_positions_value_usdc": total_bet, + "num_open_trades": open_trades, + "num_markets": len(self.positions), + "positions": self.positions, + "dry_run": self.dry_run, + } + + def cancel_all_orders(self): + """Cancel all open orders. Only works in live mode.""" + if self.dry_run or not self.clob: + print("Cannot cancel orders in dry-run mode") + return + + try: + self.clob.cancel_all() + print("All orders cancelled") + except Exception as e: + print(f"Cancellation error: {e}") diff --git a/pipeline.py b/pipeline.py new file mode 100644 index 0000000..22b6ea2 --- /dev/null +++ b/pipeline.py @@ -0,0 +1,246 @@ +#!/usr/bin/env python3 +""" +HK Weather Prediction Market Pipeline + +End-to-end pipeline: +1. Fetch weather data from Open-Meteo (WeatherNext API) and HKO +2. Extract and calibrate Hong Kong-specific forecasts +3. Scan Polymarket for relevant weather markets +4. Generate trading signals based on model edge +5. Execute trades (with dry-run safety) + +Usage: + python pipeline.py # Dry run with reporting + python pipeline.py --live # Live trading (requires keys) + python pipeline.py --schedule # Run as scheduled service +""" + +import sys +import time +import argparse +from datetime import datetime, timedelta + +from weather.hk_extractor import HKExtractor +from weather.hko_client import HKOClient +from weather.openmeteo_client import OpenMeteoClient +from markets.polymarket_client import PolymarketClient +from markets.trader import Trader +from strategy.signals import SignalGenerator +from config import POLYMARKET_PRIVATE_KEY, POLYMARKET_FUNDER + + +class Pipeline: + """Main orchestration pipeline.""" + + def __init__(self, live: bool = False, bankroll: float = 1000.0): + self.live = live + self.bankroll = bankroll + + print(f"Initializing HK Weather Prediction Market Pipeline") + print(f" Mode: {'LIVE' if live else 'DRY RUN'}") + print(f" Bankroll: ${bankroll:.2f}") + print(f" Time: {datetime.now():%Y-%m-%d %H:%M:%S %Z}") + print() + + self.hko = HKOClient() + self.openmeteo = OpenMeteoClient() + self.weather = HKExtractor() + self.polymarket = PolymarketClient() + self.signals = SignalGenerator(bankroll_usdc=bankroll) + + if live: + if not POLYMARKET_PRIVATE_KEY: + print("ERROR: POLYMARKET_PRIVATE_KEY not set in .env") + print("Falling back to dry-run mode.") + self.live = False + self.trader = Trader( + client=self.polymarket, + private_key="", + funder_address=POLYMARKET_FUNDER, + dry_run=True, + ) + else: + self.trader = Trader( + client=self.polymarket, + private_key=POLYMARKET_PRIVATE_KEY, + funder_address=POLYMARKET_FUNDER, + dry_run=False, + ) + else: + self.trader = Trader( + client=self.polymarket, + private_key="", + funder_address="", + dry_run=True, + ) + + def run(self): + """Execute a full pipeline cycle.""" + try: + self._step_fetch_data() + self._step_scan_markets() + self._step_generate_signals() + self._step_execute_trades() + self._step_report() + except KeyboardInterrupt: + print("\nPipeline interrupted.") + except Exception as e: + print(f"\nPipeline error: {e}") + import traceback + traceback.print_exc() + + def _step_fetch_data(self): + """Step 1: Fetch all weather data.""" + print("=" * 60) + print("STEP 1: Fetching weather data") + print("=" * 60) + + print(" Fetching WeatherNext forecast from Open-Meteo...") + tomorrow = (datetime.now() + timedelta(days=1)).strftime("%Y-%m-%d") + self.weather_data = self.openmeteo.get_scoring_window_summary(tomorrow) + if self.weather_data: + print(f" Temp max: {self.weather_data.get('temperature_2m_max', 'N/A')}°C") + print(f" Rain prob: {self.weather_data.get('precipitation_probability_max', 'N/A')}%") + print(f" Wind max: {self.weather_data.get('wind_speed_10m_max', 'N/A')} km/h") + else: + print(" Failed to fetch forecast data") + + print(" Fetching current HKO observations...") + self.current_obs = self.hko.get_current_weather() + if self.current_obs: + temps = self.current_obs.get("temperature", []) + if temps: + print(f" {temps[0].get('place')}: {temps[0].get('value')}°C") + warnings = self.current_obs.get("warning_message", "") + if warnings: + print(f" Warnings: {warnings}") + + print(" Fetching typhoon information...") + self.typhoon_info = self.hko.get_typhoon_info() + if self.typhoon_info: + print(f" Typhoon data: available") + + print(" Fetching HKO 9-day forecast...") + self.hko_forecast = self.hko.get_forecast() + if self.hko_forecast: + d0 = self.hko_forecast[0] + print(f" {d0['date']}: {d0['forecast_temp_min']}-{d0['forecast_temp_max']}°C, " + f"Rain: {d0.get('forecast_rain_probability', 'N/A')}") + + print() + + def _step_scan_markets(self): + """Step 2: Scan Polymarket for weather markets.""" + print("=" * 60) + print("STEP 2: Scanning Polymarket markets") + print("=" * 60) + + self.markets = self.polymarket.find_relevant_weather_markets() + + if not self.markets: + print(" No active Hong Kong weather markets found on Polymarket.") + print(" (This is expected — weather markets are less common.)") + print(" Generating standalone forecasts for potential market creation.") + else: + print(f" Found {len(self.markets)} relevant markets:") + for m in self.markets[:10]: + print(f" [{m.get('liquidity', 0):.0f} USDC liq] {m['question']}") + if m.get("outcome_prices"): + print(f" YES: {m['outcome_prices'][0]} | NO: {m['outcome_prices'][1] if len(m['outcome_prices']) > 1 else 'N/A'}") + + print() + + def _step_generate_signals(self): + """Step 3: Generate trading signals.""" + print("=" * 60) + print("STEP 3: Generating trading signals") + print("=" * 60) + + self._signals = self.signals.generate_signals() + print(self.signals.get_signal_summary()) + print() + + def _step_execute_trades(self): + """Step 4: Execute trades.""" + print("=" * 60) + print(f"STEP 4: Executing trades ({'LIVE' if self.live else 'DRY RUN'})") + print("=" * 60) + + executable = [s for s in self._signals if s.signal_type != "pass"] + + if not executable: + print(" No trades to execute (no edge above threshold)") + else: + print(f" Executing {len(executable)} trade(s)...") + for signal in executable: + result = self.trader.execute_signal(signal) + if result.success: + print(f" [{signal.signal_type}] ${result.filled_amount:.2f} - {signal.question}") + else: + print(f" FAILED: {signal.question} - {result.error}") + + summary = self.trader.get_positions_summary() + print(f"\n Positions summary:") + print(f" Total value: ${summary['total_positions_value_usdc']:.2f}") + print(f" Active trades: {summary['num_open_trades']}") + print(f" Markets: {summary['num_markets']}") + + print() + + def _step_report(self): + """Step 5: Print summary report.""" + print("=" * 60) + print("PIPELINE COMPLETE") + print("=" * 60) + print(f" Time: {datetime.now():%Y-%m-%d %H:%M:%S}") + print(f" Mode: {'LIVE' if self.live else 'DRY RUN'}") + print(f" Bankroll: ${self.bankroll:.2f}") + + positions = self.trader.get_positions_summary() + print(f" Open positions: ${positions['total_positions_value_usdc']:.2f}") + + n_active = len([s for s in self._signals if s.signal_type != "pass"]) + print(f" Active signals: {n_active}") + + print() + + +def main(): + parser = argparse.ArgumentParser( + description="HK Weather Prediction Market Pipeline", + formatter_class=argparse.RawDescriptionHelpFormatter, + epilog=""" +Examples: + python pipeline.py # Dry run with reporting + python pipeline.py --live # Live trading (requires .env keys) + python pipeline.py --schedule # Run continuously every 6 hours + python pipeline.py --bankroll 5000 # Set bankroll for Kelly sizing + """, + ) + parser.add_argument("--live", action="store_true", help="Enable live trading") + parser.add_argument("--schedule", action="store_true", help="Run on schedule (every 6 hours)") + parser.add_argument("--bankroll", type=float, default=1000.0, help="Starting bankroll in USDC") + parser.add_argument("--once", action="store_true", help="Run once and exit (default)") + + args = parser.parse_args() + + pipeline = Pipeline(live=args.live, bankroll=args.bankroll) + + if args.schedule: + print(f"Running on schedule (every 6 hours)") + print(f"Next run: {datetime.now() + timedelta(hours=6)}") + while True: + pipeline.run() + wait = 6 * 3600 + print(f"\nWaiting {wait // 3600} hours until next run...\n") + try: + time.sleep(wait) + except KeyboardInterrupt: + print("\nShutting down scheduled pipeline.") + break + else: + pipeline.run() + + +if __name__ == "__main__": + main() diff --git a/scheduled_runner.py b/scheduled_runner.py new file mode 100644 index 0000000..aa8dd02 --- /dev/null +++ b/scheduled_runner.py @@ -0,0 +1,60 @@ +#!/usr/bin/env python3 +"""Scheduled runner that executes the pipeline hourly and logs results.""" + +import time +import signal +import sys +from datetime import datetime +from pathlib import Path + +from pipeline import Pipeline + + +LOG_DIR = Path(__file__).parent / "logs" +LOG_DIR.mkdir(exist_ok=True) + +running = True + + +def signal_handler(sig, frame): + global running + print("\nShutting down scheduled runner...") + running = False + + +signal.signal(signal.SIGINT, signal_handler) +signal.signal(signal.SIGTERM, signal_handler) + + +def main(): + pipeline = Pipeline(live=False, bankroll=1000.0) + + print(f"Scheduled runner started at {datetime.now():%Y-%m-%d %H:%M:%S}") + print(f"Running every hour. Logs in {LOG_DIR}/") + print() + + while running: + try: + log_file = LOG_DIR / f"run_{datetime.now():%Y%m%d_%H%M}.txt" + with open(log_file, "w") as f: + old_stdout = sys.stdout + sys.stdout = f + pipeline.run() + sys.stdout = old_stdout + + print(f"[{datetime.now():%H:%M}] Run complete -> {log_file}") + + except Exception as e: + print(f"[{datetime.now():%H:%M}] Error: {e}") + + # Sleep until next hour + for _ in range(3600): + if not running: + break + time.sleep(1) + + print("Runner stopped.") + + +if __name__ == "__main__": + main() diff --git a/scripts/download_weights.py b/scripts/download_weights.py new file mode 100644 index 0000000..6553543 --- /dev/null +++ b/scripts/download_weights.py @@ -0,0 +1,48 @@ +#!/usr/bin/env python3 +"""Download WeatherNext model weights from Google Cloud Storage.""" + +import os +import sys +from pathlib import Path +from config import WEIGHTS_DIR, MODEL_CONFIGS + +try: + from google.cloud import storage +except ImportError: + import subprocess + subprocess.check_call([sys.executable, "-m", "pip", "install", "google-cloud-storage"]) + from google.cloud import storage + +def download_blob(bucket_name, source_blob_name, destination): + """Download a blob from GCS bucket.""" + print(f"Downloading {source_blob_name}...") + storage_client = storage.Client.create_anonymous_client() + bucket = storage_client.bucket(bucket_name) + blob = bucket.blob(source_blob_name) + blob.download_to_filename(str(destination)) + print(f" -> {destination} ({os.path.getsize(destination) / 1e6:.1f} MB)") + +def download_all(): + """Download all required model weights.""" + bucket = "dm_graphcast" + + weights = [ + "WeatherNextCyclones_Mini_<2024.npz", + "WeatherNextCyclones_Mini_<2023.npz", + ] + + for w in weights: + dst = WEIGHTS_DIR / w + if dst.exists(): + print(f"Skipping {w} (already exists)") + continue + try: + download_blob(bucket, w, dst) + except Exception as e: + print(f"Failed to download {w}: {e}") + + print("\n=== Download complete ===") + print(f"Weights stored in: {WEIGHTS_DIR}") + +if __name__ == "__main__": + download_all() diff --git a/setup.sh b/setup.sh new file mode 100644 index 0000000..d3067d7 --- /dev/null +++ b/setup.sh @@ -0,0 +1,62 @@ +#!/usr/bin/env bash +set -euo pipefail + +PROJECT_DIR="$(cd "$(dirname "$0")" && pwd)" +cd "$PROJECT_DIR" + +echo "=== Setting up HK Weather Prediction Market Pipeline ===" + +# Use Python 3.10 for JAX compatibility +export PYENV_ROOT="$HOME/.pyenv" +export PATH="$PYENV_ROOT/bin:$PATH" +eval "$(pyenv init -)" + +if ! pyenv versions | grep -q "3.10.17"; then + echo "Python 3.10.17 not found in pyenv. Please install it first." + exit 1 +fi + +if [ ! -d "venv" ]; then + pyenv local 3.10.17 + python -m venv venv +fi + +source venv/bin/activate + +echo "=== Installing dependencies ===" + +pip install --upgrade pip setuptools wheel + +# JAX with CUDA 12 support +pip install "jax[cuda12]" -f https://storage.googleapis.com/jax-releases/jax_cuda_releases.html + +# Data & scientific +pip install numpy pandas xarray scipy netCDF4 h5netcdf cfgrib +pip install zarr fsspec gcsfs + +# Weather data access +pip install ecmwf-api-client openmeteo-requests requests-cache retry-requests +pip install cdsapi + +# Polymarket +pip install py-clob-client websocket-client + +# Utilities +pip install httpx python-dotenv click rich tqdm schedule +pip install matplotlib cartopy + +# WeatherNext itself +pip install git+https://github.com/google-deepmind/weathernext.git@v0.3.0 + +echo "" +echo "=== Verifying installations ===" +python -c "import jax; print(f'JAX version: {jax.__version__}'); print(f'Devices: {jax.devices()}')" +python -c "import weathernext; print('WeatherNext imported successfully')" + +echo "" +echo "=== Setup complete ===" +echo "Next steps:" +echo " 1. source venv/bin/activate" +echo " 2. Copy .env.example to .env and fill in API keys" +echo " 3. Download model weights: python scripts/download_weights.py" +echo " 4. Run pipeline: python pipeline.py" diff --git a/strategy/__init__.py b/strategy/__init__.py new file mode 100644 index 0000000..a522877 --- /dev/null +++ b/strategy/__init__.py @@ -0,0 +1,7 @@ +"""Trading strategy components.""" + +from .calibrator import ProbabilityCalibrator +from .kelly import KellyCriterion +from .signals import SignalGenerator + +__all__ = ["ProbabilityCalibrator", "KellyCriterion", "SignalGenerator"] diff --git a/strategy/calibrator.py b/strategy/calibrator.py new file mode 100644 index 0000000..0c58124 --- /dev/null +++ b/strategy/calibrator.py @@ -0,0 +1,156 @@ +"""Probability calibration for WeatherNext forecasts. + +Converts raw model outputs into well-calibrated probabilities +suitable for prediction market trading. +""" + +import json +from datetime import datetime +from pathlib import Path +from typing import Dict, Optional, List, Tuple + +import numpy as np + + +class ProbabilityCalibrator: + """ + Calibrate raw model probabilities using historical performance. + + Methods: + - Platt scaling (logistic regression on historical outcomes) + - Isotonic regression (non-parametric) + - Ensemble (combine multiple calibration methods) + """ + + def __init__(self, calibration_file: str = "data/calibration_history.json"): + self.calibration_file = Path(__file__).parent.parent / calibration_file + self.history: List[Dict] = self._load_history() + self.platt_params: Dict[str, Tuple[float, float]] = {} + self._fit() + + def _load_history(self) -> List[Dict]: + """Load historical forecast vs outcome data.""" + if self.calibration_file.exists(): + try: + with open(self.calibration_file) as f: + return json.load(f) + except Exception: + return [] + return [] + + def save_history(self): + """Save calibration history.""" + self.calibration_file.parent.mkdir(parents=True, exist_ok=True) + with open(self.calibration_file, "w") as f: + json.dump(self.history, f, indent=2) + + def record_outcome( + self, + date: str, + variable: str, + predicted_probability: float, + actual_outcome: bool, + ): + """Record a prediction-outcome pair for future calibration.""" + self.history.append({ + "date": date, + "variable": variable, + "predicted_probability": predicted_probability, + "actual_outcome": actual_outcome, + "recorded_at": datetime.now().isoformat(), + }) + self.save_history() + self._fit() # Re-fit on new data + + def _fit(self): + """Fit Platt scaling parameters from history.""" + by_variable: Dict[str, List[Tuple[float, int]]] = {} + for record in self.history: + var = record["variable"] + if var not in by_variable: + by_variable[var] = [] + by_variable[var].append(( + record["predicted_probability"] / 100.0, + 1 if record["actual_outcome"] else 0, + )) + + for var, data in by_variable.items(): + if len(data) >= 5: + try: + from sklearn.linear_model import LogisticRegression + X = np.array([[d[0]] for d in data]) + y = np.array([d[1] for d in data]) + lr = LogisticRegression() + lr.fit(X, y) + self.platt_params[var] = (lr.coef_[0][0], lr.intercept_[0]) + except ImportError: + # Fallback: simple linear correction + self._simple_fit(var, data) + except Exception: + self._simple_fit(var, data) + elif len(data) >= 2: + self._simple_fit(var, data) + + def _simple_fit(self, var: str, data: List[Tuple[float, int]]): + """Simple linear calibration for small datasets.""" + probs = np.array([d[0] for d in data]) + outcomes = np.array([d[1] for d in data]) + mean_prob = probs.mean() + mean_outcome = outcomes.mean() + + slope = 1.0 + intercept = mean_outcome - mean_prob + + self.platt_params[var] = (slope, intercept) + + def calibrate(self, variable: str, raw_probability: float) -> float: + """Calibrate a raw probability (0-100) to a calibrated one.""" + x = raw_probability / 100.0 + + if variable in self.platt_params: + a, b = self.platt_params[variable] + calibrated = 1.0 / (1.0 + np.exp(-(a * x + b))) + return float(np.clip(calibrated * 100.0, 0.5, 99.5)) + + return float(np.clip(raw_probability, 0.5, 99.5)) + + def ensemble_calibrate( + self, variable: str, raw_probability: float + ) -> Tuple[float, float]: + """ + Return (calibrated_probability, confidence_interval_width). + Confidence width shrinks with more historical data. + """ + cal_prob = self.calibrate(variable, raw_probability) + + n_obs = sum(1 for h in self.history if h["variable"] == variable) + if n_obs < 5: + ci_width = 15.0 + elif n_obs < 20: + ci_width = 10.0 + elif n_obs < 50: + ci_width = 5.0 + else: + ci_width = 3.0 + + return cal_prob, ci_width + + def get_calibration_stats(self, variable: str) -> Dict: + """Get calibration statistics for a variable.""" + relevant = [h for h in self.history if h["variable"] == variable] + if not relevant: + return {"n_observations": 0, "brier_score": None, "calibration_error": None} + + preds = np.array([h["predicted_probability"] / 100.0 for h in relevant]) + outcomes = np.array([1 if h["actual_outcome"] else 0 for h in relevant]) + + brier = float(np.mean((preds - outcomes) ** 2)) + cal_error = float(np.abs(preds.mean() - outcomes.mean())) + + return { + "n_observations": len(relevant), + "brier_score": brier, + "calibration_error": cal_error, + "mean_prediction": float(preds.mean() * 100), + "mean_outcome": float(outcomes.mean() * 100), + } diff --git a/strategy/kelly.py b/strategy/kelly.py new file mode 100644 index 0000000..b74be57 --- /dev/null +++ b/strategy/kelly.py @@ -0,0 +1,181 @@ +"""Kelly Criterion position sizing for prediction market betting. + +Implements fractional Kelly to control risk while maximizing +log-wealth growth based on model edge vs market implied probability. +""" + +import numpy as np +from dataclasses import dataclass + +from config import KELLY_FRACTION, MAX_POSITION_USDC + + +@dataclass +class KellyResult: + """Result of Kelly sizing calculation.""" + full_kelly_fraction: float # Fraction of bankroll to bet (full Kelly) + fractional_kelly: float # Fraction after applying Kelly fraction + size_usdc: float # Absolute bet size in USDC + kelly_active: bool # Whether full Kelly recommends a bet + edge: float # Edge in decimal (not bps) + log_utility: float # Expected log-utility gain + + +class KellyCriterion: + """ + Kelly Criterion for binary prediction markets. + + For binary markets: + f* = p - q / (b) + where: + p = our estimated probability of winning + q = 1 - p + b = net odds received (payout / bet - 1) + + In prediction markets: + If we buy YES at price P, and it resolves YES, we get 1. + So b = (1-P)/P if buying YES, or P/(1-P) if buying NO. + """ + + def __init__(self, bankroll_usdc: float = 1000.0, fraction: float = KELLY_FRACTION): + self.bankroll = bankroll_usdc + self.fraction = fraction + + def size_bet( + self, + our_probability: float, # Our probability (0-100 or 0-1) + market_probability: float, # Market probability (0-100 or 0-1) + side: str = "buy_yes", # "buy_yes" or "buy_no" + max_size: float = MAX_POSITION_USDC, + ) -> KellyResult: + """ + Calculate Kelly-optimal bet size. + + our_probability: Our model's probability of YES outcome (0-1 or 0-100) + market_probability: Market-implied probability of YES outcome (0-1 or 0-100) + """ + # Normalize to 0-1 range + if our_probability > 1: + our_probability /= 100.0 + if market_probability > 1: + market_probability /= 100.0 + + # Clamp to avoid division by zero or log(0) + our_probability = np.clip(our_probability, 0.001, 0.999) + market_probability = np.clip(market_probability, 0.001, 0.999) + + if side == "buy_yes": + # Buy YES: we win 1-P per share at cost P + b = (1.0 - market_probability) / market_probability # Net odds + p = our_probability + q = 1.0 - our_probability + else: + # Buy NO: symmetric + b = market_probability / (1.0 - market_probability) + p = 1.0 - our_probability # We win if NO + q = our_probability + + # Kelly formula: f* = (p * b - q) / b = p - q/b + if b > 0: + full_kelly = p - q / b + else: + full_kelly = 0.0 + + # Edge in decimal + if side == "buy_yes": + edge = our_probability - market_probability + else: + edge = (1.0 - our_probability) - (1.0 - market_probability) + edge = market_probability - our_probability # Same thing + + # Only bet when we have positive edge + kelly_active = full_kelly > 0.001 + + if not kelly_active: + return KellyResult( + full_kelly_fraction=0.0, + fractional_kelly=0.0, + size_usdc=0.0, + kelly_active=False, + edge=edge, + log_utility=0.0, + ) + + # Apply fraction for safety + fractional_kelly = full_kelly * self.fraction + size_usdc = min(fractional_kelly * self.bankroll, max_size) + + # Log utility uses fractions of bankroll + log_utility = self._expected_log_utility( + p_win=our_probability, + market_price=market_probability, + side=side, + bet_fraction=min(fractional_kelly, 0.99) if kelly_active else 0.0, + ) + + return KellyResult( + full_kelly_fraction=full_kelly, + fractional_kelly=fractional_kelly, + size_usdc=size_usdc, + kelly_active=kelly_active, + edge=edge, + log_utility=log_utility, + ) + + def compare_sides( + self, + our_probability: float, + market_probability: float, + ) -> dict: + """Compare betting YES vs NO and return the better side.""" + yes_result = self.size_bet(our_probability, market_probability, "buy_yes") + no_result = self.size_bet(our_probability, market_probability, "buy_no") + + if yes_result.size_usdc > no_result.size_usdc: + return { + "recommended_side": "buy_yes", + "size_usdc": yes_result.size_usdc, + "edge": yes_result.edge, + "log_utility": yes_result.log_utility, + } + else: + return { + "recommended_side": "buy_no", + "size_usdc": no_result.size_usdc, + "edge": no_result.edge, + "log_utility": no_result.log_utility, + } + + def update_bankroll(self, new_bankroll: float): + """Update bankroll after wins/losses.""" + self.bankroll = new_bankroll + + @staticmethod + def _expected_log_utility( + p_win: float, + market_price: float, + side: str, + bet_fraction: float, + ) -> float: + """Calculate expected log-utility (Kelly criterion) of a fractional bet.""" + if bet_fraction <= 0: + return 0.0 + + if side == "buy_yes": + win_mult = (1.0 - market_price) / market_price + else: + win_mult = market_price / (1.0 - market_price) + + # bet_fraction is fraction of bankroll + # Win: bankroll becomes bankroll * (1 + bet_fraction * win_mult) + # Lose: bankroll becomes bankroll * (1 - bet_fraction) + total_after_win = 1.0 + bet_fraction * win_mult + total_after_loss = 1.0 - bet_fraction + + if total_after_loss <= 0: + return -999.0 + + if side == "buy_yes": + return p_win * np.log(max(1e-10, total_after_win)) + (1.0 - p_win) * np.log(max(1e-10, total_after_loss)) + else: + return (1.0 - p_win) * np.log(max(1e-10, total_after_win)) + p_win * np.log(max(1e-10, total_after_loss)) diff --git a/strategy/signals.py b/strategy/signals.py new file mode 100644 index 0000000..7ff665c --- /dev/null +++ b/strategy/signals.py @@ -0,0 +1,206 @@ +"""Signal generator for HK weather prediction market trading. + +Combines model forecasts, probability calibration, and Kelly sizing +to generate trading signals for Polymarket execution. +""" + +from datetime import datetime, timedelta +from typing import Optional, Dict, List + +from weather.hk_extractor import HKExtractor +from strategy.calibrator import ProbabilityCalibrator +from strategy.kelly import KellyCriterion, KellyResult +from markets.polymarket_client import PolymarketClient +from markets.trader import TradeSignal +from config import MIN_EDGE_BPS + + +class SignalGenerator: + """Generate trading signals from weather forecasts and market prices.""" + + def __init__( + self, + bankroll_usdc: float = 1000.0, + min_edge_bps: float = MIN_EDGE_BPS, + ): + self.weather = HKExtractor() + self.polymarket = PolymarketClient() + self.calibrator = ProbabilityCalibrator() + self.kelly = KellyCriterion(bankroll_usdc=bankroll_usdc) + self.min_edge_bps = min_edge_bps + self.signals: List[TradeSignal] = [] + + def generate_signals(self) -> List[TradeSignal]: + """Generate all trading signals for available markets.""" + self.signals = [] + + markets = self.polymarket.find_relevant_weather_markets() + + if not markets: + print("No relevant markets found on Polymarket") + self._generate_standalone_signals() + return self.signals + + for market in markets: + signal = self._analyze_market(market) + if signal: + self.signals.append(signal) + + self.signals.sort(key=lambda s: abs(s.edge_bps), reverse=True) + return self.signals + + def _analyze_market(self, market: Dict) -> Optional[TradeSignal]: + """Analyze a single market and generate a signal.""" + condition_id = market["condition_id"] + question = market["question"].lower() + + if not market.get("active") or market.get("closed"): + return None + + if market.get("liquidity", 0) < 50: + return None # Too illiquid + + # Get market-implied probability + market_prob = self.polymarket.get_market_implied_probability(condition_id, 0) + if market_prob is None: + return None + + # Determine what we're predicting + model_prob, variable = self._get_model_probability(question) + + if model_prob is None: + return None + + # Calibrate our probability + cal_prob = self.calibrator.calibrate(variable, model_prob) + + # Calculate edge + edge_bps = (cal_prob - market_prob) * 100 # Convert to basis points + + if abs(edge_bps) < self.min_edge_bps: + return TradeSignal( + market_id=market["id"], + condition_id=condition_id, + question=question, + outcome_index=0, + outcome_label=market["outcomes"][0] if market.get("outcomes") else "Yes", + model_probability=cal_prob, + market_probability=market_prob, + edge_bps=edge_bps, + recommended_size_usdc=0, + max_size_usdc=0, + signal_type="pass", + ) + + # Determine side + side = "buy_yes" if edge_bps > 0 else "buy_no" + side_prob = cal_prob if side == "buy_yes" else 100 - cal_prob + + # Kelly sizing + kelly_result = self.kelly.size_bet( + our_probability=cal_prob, + market_probability=market_prob, + side=side, + ) + + return TradeSignal( + market_id=market["id"], + condition_id=condition_id, + question=question, + outcome_index=0, + outcome_label=market["outcomes"][0] if market.get("outcomes") else "Yes", + model_probability=cal_prob, + market_probability=market_prob, + edge_bps=edge_bps, + recommended_size_usdc=kelly_result.size_usdc, + max_size_usdc=kelly_result.size_usdc, + signal_type=side, + ) + + def _get_model_probability(self, question: str) -> tuple: + """Get our model's probability for a given market question.""" + question = question.lower() + + if "rain" in question or "precipitation" in question: + prob = self.weather.should_bet_rain_tomorrow() + return (prob, "rain") if prob is not None else (None, "") + + if "temperature" in question and ("above" in question or "exceed" in question): + if "30" in question or "thirty" in question: + prob = self.weather.should_bet_temp_above(30.0) + elif "35" in question or "thirty five" in question: + prob = self.weather.should_bet_temp_above(35.0) + else: + prob = self.weather.should_bet_temp_above(33.0) + return (prob, "temperature") if prob is not None else (None, "") + + if "typhoon" in question or "t8" in question or "tropical cyclone" in question: + forecast = self.weather.get_hk_forecast() + typhoon = forecast.get("typhoon_info", {}) + prob = 30.0 if typhoon else 5.0 + return (prob, "typhoon") + + if "weather" in question or "storm" in question: + forecast = self.weather.get_combined_tomorrow_forecast() + tomorrow = forecast.get("tomorrow", {}) + if tomorrow: + rain_prob = tomorrow.get("precipitation_probability_calibrated", 50) + return (rain_prob, "rain") + + return (None, "") + + def _generate_standalone_signals(self): + """Generate signals even when no Polymarket markets exist. + Useful for tracking model predictions and for creating new markets. + """ + tomorrow = datetime.now() + timedelta(days=1) + forecast = self.weather.get_hk_forecast() + consensus = forecast.get("consensus", {}) + tmrw = consensus.get("tomorrow", {}) + + if tmrw: + self.signals.append(TradeSignal( + market_id="standalone", + condition_id="standalone", + question=f"Will it rain in Hong Kong on {tomorrow:%Y-%m-%d}?", + outcome_index=0, + outcome_label="Yes", + model_probability=tmrw.get("precipitation_probability_calibrated", 50), + market_probability=50.0, + edge_bps=0, + recommended_size_usdc=0, + max_size_usdc=0, + signal_type="pass", + )) + + print(f"\nGenerated {len(self.signals)} standalone signals") + for s in self.signals: + print(f" {s.question} -> P={s.model_probability:.1f}%") + + def get_signal_summary(self) -> str: + """Get a human-readable summary of current signals.""" + if not self.signals: + return "No signals generated." + + lines = [] + active = [s for s in self.signals if s.signal_type != "pass"] + passed = [s for s in self.signals if s.signal_type == "pass"] + + lines.append(f"\n=== Signal Summary ({datetime.now():%Y-%m-%d %H:%M}) ===") + lines.append(f"Active signals: {len(active)}") + lines.append(f"Passed (no edge): {len(passed)}") + lines.append("") + + if active: + lines.append("TRADE SIGNALS:") + for s in active: + lines.append(f" [{s.signal_type.upper()}] {s.question}") + lines.append(f" Model: {s.model_probability:.1f}% | Market: {s.market_probability:.1f}%") + lines.append(f" Edge: {s.edge_bps:.0f}bps | Size: ${s.recommended_size_usdc:.2f}") + + if passed: + lines.append("PASSED (edge < threshold):") + for s in passed[:5]: # Limit to 5 + lines.append(f" {s.question} (edge: {s.edge_bps:.0f}bps)") + + return "\n".join(lines) diff --git a/weather/__init__.py b/weather/__init__.py new file mode 100644 index 0000000..dfa3198 --- /dev/null +++ b/weather/__init__.py @@ -0,0 +1,7 @@ +"""Weather data for Hong Kong weather prediction markets.""" + +from .openmeteo_client import OpenMeteoClient +from .hko_client import HKOClient +from .hk_extractor import HKExtractor + +__all__ = ["OpenMeteoClient", "HKOClient", "HKExtractor"] diff --git a/weather/hk_extractor.py b/weather/hk_extractor.py new file mode 100644 index 0000000..4dfc742 --- /dev/null +++ b/weather/hk_extractor.py @@ -0,0 +1,197 @@ +"""Extract Hong Kong specific forecasts from global model outputs. + +Handles regional extraction, downscaling hints, and local calibration +based on HKO station data for the Hong Kong region. +""" + +from datetime import datetime, timedelta +from typing import Optional, Dict, List + +import numpy as np +import pandas as pd + +from config import HK_COORDS, HK_BBOX +from .openmeteo_client import OpenMeteoClient +from .hko_client import HKOClient + + +class HKExtractor: + """Extract and calibrate HK-specific weather forecasts from global models.""" + + # Known stations for calibration (HKO stations with good historical data) + CALIBRATION_STATIONS = [ + "Hong Kong Observatory", # Tsim Sha Tsui + "Chek Lap Kok", # Airport + "Sha Tin", + "Tuen Mun", + "Sai Kung", + "Ta Kwu Ling", + "Sheung Shui", + "Stanley", + ] + + # Calibration offsets - will be learned over time + # (model_bias, model_std) for key variables + DEFAULT_BIAS = { + "temperature_2m_max": 0.0, + "temperature_2m_min": 0.0, + "precipitation_probability_max": 0.0, + "wind_speed_10m_max": 0.0, + } + + def __init__(self, calibrate: bool = True): + self.openmeteo = OpenMeteoClient() + self.hko = HKOClient() + self.calibrate = calibrate + self.bias_model = self.DEFAULT_BIAS.copy() + self._load_calibration() + + def _load_calibration(self): + """Load calibration params from stored file if available.""" + import os + import json + path = os.path.join(os.path.dirname(__file__), "..", "data", "calibration.json") + if os.path.exists(path): + try: + with open(path) as f: + stored = json.load(f) + self.bias_model.update(stored.get("bias", {})) + except Exception: + pass + + def save_calibration(self): + """Save calibration params for future runs.""" + import os + import json + path = os.path.join(os.path.dirname(__file__), "..", "data", "calibration.json") + with open(path, "w") as f: + json.dump({"bias": self.bias_model, "updated": datetime.now().isoformat()}, f, indent=2) + + def get_hk_forecast(self, lead_days: int = 7) -> Dict: + """Get calibrated HK-specific forecast combining multiple sources.""" + forecast = { + "fetch_time": datetime.now().isoformat(), + "sources": {}, + } + + wnext = self.openmeteo.get_forecast(lead_days=lead_days) + if wnext is not None: + forecast["sources"]["weathernext"] = self._calibrate_forecast(wnext) + + hko_fc = self.hko.get_forecast() + if hko_fc: + forecast["sources"]["hko"] = hko_fc + + current = self.hko.get_current_weather() + if current: + forecast["current_observations"] = current + + typhoon = self.hko.get_typhoon_info() + if typhoon: + forecast["typhoon_info"] = typhoon + + forecast["consensus"] = self._build_consensus(forecast) + + return forecast + + def _calibrate_forecast(self, df: pd.DataFrame) -> List[Dict]: + """Apply calibration to model forecast and return structured data.""" + results = [] + for idx, row in df.iterrows(): + day = { + "date": idx.strftime("%Y-%m-%d"), + "temp_max_calibrated": float(row.get("temperature_2m_max", np.nan)) + self.bias_model.get("temperature_2m_max", 0), + "temp_min_calibrated": float(row.get("temperature_2m_min", np.nan)) + self.bias_model.get("temperature_2m_min", 0), + "temp_max_raw": float(row.get("temperature_2m_max", np.nan)), + "temp_min_raw": float(row.get("temperature_2m_min", np.nan)), + "precipitation_probability_calibrated": min(100, max(0, float(row.get("precipitation_probability_max", 0)) + self.bias_model.get("precipitation_probability_max", 0))), + "precipitation_probability_raw": float(row.get("precipitation_probability_max", 0)), + "precipitation_sum": float(row.get("precipitation_sum", 0)), + "wind_speed_max_calibrated": float(row.get("wind_speed_10m_max", np.nan)) + self.bias_model.get("wind_speed_10m_max", 0), + "wind_speed_max_raw": float(row.get("wind_speed_10m_max", np.nan)), + "wind_gusts_max": float(row.get("wind_gusts_10m_max", np.nan)), + } + results.append(day) + return results + + def _build_consensus(self, forecast: Dict) -> Dict: + """Build a consensus forecast from all available sources.""" + consensus = {} + + if "weathernext" in forecast.get("sources", {}): + w = forecast["sources"]["weathernext"] + if w: + d0 = w[0] + consensus["tomorrow"] = d0 + + if "hko" in forecast.get("sources", {}): + h = forecast["sources"]["hko"] + if h and len(h) > 0: + consensus["hko_tomorrow"] = h[0] + + current = forecast.get("current_observations", {}) + if current: + consensus["current_temp"] = ( + current.get("temperature", [{}])[0].get("value") if current.get("temperature") else None + ) + + return consensus + + def get_combined_tomorrow_forecast(self) -> Dict: + """Get a single combined forecast for 'tomorrow' from all sources.""" + fc = self.get_hk_forecast() + return fc.get("consensus", {}) + + def should_bet_rain_tomorrow(self) -> Optional[float]: + """Returns model-implied probability of rain tomorrow (0-100).""" + fc = self.get_hk_forecast() + consensus = fc.get("consensus", {}) + tomorrow = consensus.get("tomorrow", {}) + hko = consensus.get("hko_tomorrow", {}) + + probs = [] + if "precipitation_probability_calibrated" in tomorrow: + probs.append(tomorrow["precipitation_probability_calibrated"]) + + hko_prob_str = hko.get("forecast_rain_probability", "") + if hko_prob_str: + try: + nums = [int(x.replace("%", "")) for x in hko_prob_str.split("/")] + probs.append(max(nums)) + except (ValueError, AttributeError): + pass + + if not probs: + return None + return float(np.mean(probs)) + + def should_bet_temp_above(self, threshold: float = 30.0) -> Optional[float]: + """Returns model-implied probability that temp exceeds threshold tomorrow.""" + fc = self.get_hk_forecast() + consensus = fc.get("consensus", {}) + tomorrow = consensus.get("tomorrow", {}) + hko = consensus.get("hko_tomorrow", {}) + + temp_max_raw = tomorrow.get("temp_max_raw", np.nan) + temp_max_cal = tomorrow.get("temp_max_calibrated", np.nan) + + # Simple: if calibrated max is above threshold, probability from how far above + if not np.isnan(temp_max_cal): + excess = temp_max_cal - threshold + prob = min(100, max(0, 50 + excess * 20)) # Simple sigmoid-like + return prob + + return None + + def update_calibration(self, forecast_date: str, observed: Dict): + """Update calibration based on observed vs predicted.""" + # This would be called after the scoring window closes + # Simple exponential moving average of errors + alpha = 0.1 + for var in self.bias_model: + if var in observed and var.replace("_calibrated", "_raw") in observed: + # We'd need to store the forecast that was made for this date + # This is a placeholder for the calibration loop + pass + + self.save_calibration() diff --git a/weather/hko_client.py b/weather/hko_client.py new file mode 100644 index 0000000..86a23f0 --- /dev/null +++ b/weather/hko_client.py @@ -0,0 +1,177 @@ +"""Hong Kong Observatory data client. + +Fetches real-time weather observations, warnings, and forecasts from HKO. +Uses public RSS/JSON feeds where available. +""" + +import json +from datetime import datetime +from typing import Optional, List + +import requests + + +class HKOClient: + """Fetch data from Hong Kong Observatory's public data feeds.""" + + BASE_URL = "https://data.weather.gov.hk" + + ENDPOINTS = { + "current_weather": f"{BASE_URL}/weatherAPI/opendata/weather.php?dataType=rhrread&lang=en", + "9day_forecast": f"{BASE_URL}/weatherAPI/opendata/weather.php?dataType=fnd&lang=en", + "warnings": f"{BASE_URL}/weatherAPI/opendata/weather.php?dataType=warnsum&lang=en", + "local_forecast": f"{BASE_URL}/weatherAPI/opendata/weather.php?dataType=flw&lang=en", + "swt": f"{BASE_URL}/weatherAPI/opendata/weather.php?dataType=swt&lang=en", + "rainfall": f"{BASE_URL}/weatherAPI/opendata/weather.php?dataType=rfmap", + "lightning": f"{BASE_URL}/weatherAPI/opendata/weather.php?dataType=lmap", + "uv_index": f"{BASE_URL}/weatherAPI/opendata/weather.php?dataType=uvi", + "earthquake": f"{BASE_URL}/weatherAPI/opendata/weather.php?dataType=eq", + } + + def __init__(self): + self.session = requests.Session() + self.session.headers.update({ + "User-Agent": "HK-Weather-Market/1.0", + "Accept": "application/json", + }) + + def get_current_weather(self) -> Optional[dict]: + """Get current weather observations for Hong Kong.""" + try: + resp = self.session.get(self.ENDPOINTS["current_weather"], timeout=15) + resp.raise_for_status() + data = resp.json() + + result = { + "update_time": data.get("updateTime", ""), + "temperature": [], + "humidity": [], + "rainfall": [], + "warning_message": data.get("warningMessage", ""), + "rainstorm_reminder": data.get("rainstormReminder", ""), + "mintemp_from00to09": data.get("mintempFrom00To09", ""), + } + + # Parse temperature data + for record in data.get("temperature", {}).get("data", []): + result["temperature"].append({ + "place": record.get("place"), + "value": record.get("value"), + "unit": record.get("unit", "C"), + }) + + for record in data.get("humidity", {}).get("data", []): + result["humidity"].append({ + "place": record.get("place"), + "value": record.get("value"), + "unit": record.get("unit", "%"), + }) + + # Parse rainfall + for record in data.get("rainfall", {}).get("data", []): + result["rainfall"].append({ + "place": record.get("place"), + "max": record.get("max"), + "min": record.get("min"), + "unit": record.get("unit", "mm"), + }) + + return result + + except Exception as e: + print(f"HKO current weather error: {e}") + return None + + def get_forecast(self) -> Optional[List[dict]]: + """Get 9-day weather forecast from HKO.""" + try: + resp = self.session.get(self.ENDPOINTS["9day_forecast"], timeout=15) + resp.raise_for_status() + data = resp.json() + + forecasts = [] + for day in data.get("weatherForecast", []): + forecasts.append({ + "date": day.get("forecastDate"), + "week": day.get("week"), + "forecast_wind": day.get("forecastWind"), + "forecast_weather": day.get("forecastWeather"), + "forecast_temp_max": day.get("forecastMaxtemp", {}).get("value"), + "forecast_temp_min": day.get("forecastMintemp", {}).get("value"), + "forecast_humidity_max": day.get("forecastMaxrh", {}).get("value"), + "forecast_humidity_min": day.get("forecastMinrh", {}).get("value"), + "forecast_rain_probability": self._parse_probability( + day.get("forecastRainProbability", {}) + ), + "psr": day.get("PSR"), + }) + + return forecasts + + except Exception as e: + print(f"HKO forecast error: {e}") + return None + + def get_warnings(self) -> Optional[dict]: + """Get current weather warnings and signals.""" + try: + resp = self.session.get(self.ENDPOINTS["warnings"], timeout=15) + resp.raise_for_status() + return resp.json() + except Exception as e: + print(f"HKO warnings error: {e}") + return None + + def get_typhoon_info(self) -> Optional[dict]: + """Get tropical cyclone warnings and tracks from HKO.""" + try: + resp = self.session.get(self.ENDPOINTS["swt"], timeout=15) + resp.raise_for_status() + return resp.json() + except Exception as e: + print(f"HKO typhoon info error: {e}") + return None + + def get_rainfall_map(self) -> Optional[dict]: + """Get recent rainfall distribution.""" + try: + resp = self.session.get(self.ENDPOINTS["rainfall"], timeout=15) + resp.raise_for_status() + return resp.json() + except Exception as e: + print(f"HKO rainfall map error: {e}") + return None + + def is_t8_active(self) -> bool: + """Check if T8 or higher typhoon signal is active.""" + warnings = self.get_warnings() + if not warnings: + return False + + for detail in warnings.get("details", []): + name = detail.get("name", "") + if any(signal in name for signal in ["8", "9", "10"]): + if "tropical cyclone" in name.lower() or "typhoon" in name.lower(): + return True + return False + + def get_current_signal_level(self) -> int: + """Get current typhoon signal level (0, 1, 3, 8, 9, 10).""" + typhoon = self.get_typhoon_info() + if not typhoon: + return 0 + + for signal_text in ["T1", "T3", "T8", "T9", "T10"]: + if signal_text in str(typhoon): + return int(signal_text[1:]) + return 0 + + @staticmethod + def _parse_probability(prob_data: dict) -> Optional[str]: + """Parse HKO rainfall probability data.""" + values = [] + for period in prob_data.get("probability", []): + perc = period.get("value") + if perc: + values.append(perc) + return " / ".join(values) if values else None diff --git a/weather/model_runner.py b/weather/model_runner.py new file mode 100644 index 0000000..3674e1a --- /dev/null +++ b/weather/model_runner.py @@ -0,0 +1,216 @@ +"""Local WeatherNext model runner using JAX/GPU. + +Runs WeatherNext Cyclones Mini on RTX 4070 SUPER (12GB VRAM). +""" + +import os +from datetime import datetime +from pathlib import Path +from typing import Optional, Dict + +import numpy as np +import pandas as pd + +from config import WEIGHTS_DIR, HK_BBOX, HK_COORDS + + +class WeatherNextRunner: + """Run WeatherNext model locally for Hong Kong region forecasts.""" + + def __init__(self, model_name: str = "WeatherNextCyclones_Mini_<2024"): + self.model_name = model_name + self.weights_path = WEIGHTS_DIR / f"{model_name}.npz" + self.model = None + self.params = None + self.state = None + self.initialized = False + + def check_ready(self) -> bool: + """Check if model weights exist and JAX is working.""" + if not self.weights_path.exists(): + print(f"Weights not found: {self.weights_path}") + print("Run: python scripts/download_weights.py") + return False + + try: + import jax + devices = jax.devices() + if not devices: + print("No JAX devices found") + return False + print(f"JAX devices: {devices}") + return True + except ImportError: + print("JAX not installed") + return False + + def load_model(self) -> bool: + """Load WeatherNext model and weights.""" + if not self.check_ready(): + return False + + try: + import jax + import jax.numpy as jnp + from weathernext.models.fgn import FGN, FGNConfig + + print(f"Loading model weights from {self.weights_path}...") + params = dict(np.load(self.weights_path, allow_pickle=True)) + + config = FGNConfig( + resolution=60, # 1° for Mini + num_layers=12, + model_dim=384, + num_heads=6, + mlp_ratio=4.0, + patch_size=4, + max_path_length=240, + ) + + model = FGN(config) + self.model = model + self.params = params + self.initialized = True + print("Model loaded successfully") + + return True + + except Exception as e: + print(f"Failed to load model: {e}") + print("Using Open-Meteo WeatherNext API as fallback.") + return False + + def load_initial_state(self, source: str = "era5") -> Optional[Dict]: + """Load initial atmospheric state for model input. + + source: 'era5' (reanalysis), 'hres' (operational), or 'gfs' + """ + import xarray as xr + + try: + if source == "era5": + ds = self._load_era5_latest() + elif source == "gfs": + ds = self._load_gfs_latest() + else: + print(f"Unknown source: {source}") + return None + + return self._preprocess_for_model(ds) + + except Exception as e: + print(f"Failed to load initial state: {e}") + return None + + def run_forecast(self, lead_hours: int = 120, steps: int = 20) -> Optional[pd.DataFrame]: + """Run an autoregressive forecast for Hong Kong region. + + lead_hours: Total forecast hours + steps: Number of model steps (at 6h per step for Mini) + """ + if not self.initialized: + if not self.load_model(): + return None + + import jax + import jax.numpy as jnp + + initial_state = self.load_initial_state("era5") + if initial_state is None: + return None + + try: + input_tensor = jnp.array(initial_state["fields"]) + + @jax.jit + def step_fn(params, state): + return self.model.apply({"params": params}, state) + + results = [] + current = input_tensor + for i in range(steps): + current = step_fn(self.params, current) + if (i + 1) * 6 <= lead_hours: + results.append(self._extract_hk_region(np.array(current), i * 6 + 6)) + + return self._format_forecast_df(results) + + except Exception as e: + print(f"Model inference failed: {e}") + print("Falling back to API-based forecasts.") + return None + + def _extract_hk_region(self, field: np.ndarray, lead_hour: int) -> Dict: + """Extract Hong Kong region data from global field. + + With 1° resolution, HK is roughly a single grid cell. + """ + lat_idx = slice( + int((HK_BBOX["lat_min"] + 90) / 1.0), + int((HK_BBOX["lat_max"] + 90) / 1.0) + 1, + ) + lon_idx = slice( + int((HK_BBOX["lon_min"] + 180) / 1.0), + int((HK_BBOX["lon_max"] + 180) / 1.0) + 1, + ) + + # Placeholder - actual variable mapping depends on WeatherNext output channels + hk_slice = field[..., lat_idx, lon_idx] + return { + "lead_hour": lead_hour, + "temperature_2m_mean": float(np.mean(hk_slice[0]) if hk_slice.size > 0 else np.nan), + "precipitation_mean": float(np.mean(hk_slice[-2]) if hk_slice.size > 1 else np.nan), + } + + def _format_forecast_df(self, results: list) -> pd.DataFrame: + """Format forecast results as DataFrame.""" + if not results: + return pd.DataFrame() + + df = pd.DataFrame(results) + now = datetime.now() + df["valid_time"] = [now + pd.Timedelta(hours=r["lead_hour"]) for r in results] + return df.set_index("valid_time") + + def _load_era5_latest(self): + """Load latest ERA5 data. Requires CDS API setup.""" + import xarray as xr + from datetime import datetime, timedelta + + cds_ds = xr.open_dataset( + "https://storage.googleapis.com/dm_graphcast/dataset/dataset_test.nc", + engine="h5netcdf", + ) + return cds_ds + + def _load_gfs_latest(self): + """Load latest GFS analysis.""" + import xarray as xr + from datetime import datetime + + now = datetime.utcnow() + url = f"https://nomads.ncep.noaa.gov/dods/gfs_0p25/gfs{now:%Y%m%d}/gfs_0p25_00z" + try: + return xr.open_dataset(url) + except Exception: + return None + + def _preprocess_for_model(self, ds) -> Dict: + """Convert raw dataset to model input format.""" + required_vars = [ + "2m_temperature", "10m_u_component_of_wind", "10m_v_component_of_wind", + "mean_sea_level_pressure", "geopotential", "specific_humidity", + "temperature", "u_component_of_wind", "v_component_of_wind", + ] + + fields = [] + for var in required_vars: + if var in ds: + fields.append(np.array(ds[var].isel(time=-1))) + else: + fields.append(np.zeros((721, 1440))) + + return { + "fields": np.stack(fields), + "timestamp": str(ds.time.isel(time=-1).values), + } diff --git a/weather/openmeteo_client.py b/weather/openmeteo_client.py new file mode 100644 index 0000000..219e6ff --- /dev/null +++ b/weather/openmeteo_client.py @@ -0,0 +1,250 @@ +"""Open-Meteo client for WeatherNext and other weather model data.""" + +from datetime import datetime, timedelta +from typing import Optional + +import numpy as np +import pandas as pd + +try: + import openmeteo_requests + import requests_cache + from retry_requests import retry +except ImportError: + import subprocess, sys + subprocess.check_call([sys.executable, "-m", "pip", "install", "openmeteo-requests", "requests-cache", "retry-requests"]) + import openmeteo_requests + import requests_cache + from retry_requests import retry + +from config import HK_COORDS, HK_BBOX, OPENMETEO_API_KEY + + +class OpenMeteoClient: + """Fetch weather data from Open-Meteo including WeatherNext model outputs.""" + + BASE_URL = "https://api.open-meteo.com/v1/" + + # Available weather models via Open-Meteo + MODELS = { + "weathernext": "google_weathernext", + "ecmwf": "ecmwf_ifs04", + "gfs": "gfs_seamless", + } + + def __init__(self, model: str = "weathernext", cache_ttl: int = 3600): + cache = requests_cache.CachedSession('.openmeteo_cache', expire_after=cache_ttl) + retry_session = retry(cache, retries=3, backoff_factor=0.2) + self.client = openmeteo_requests.Client(session=retry_session) + self.model = model + self.params = { + "latitude": HK_COORDS["hko_headquarters"][0], + "longitude": HK_COORDS["hko_headquarters"][1], + "timezone": "Asia/Hong_Kong", + } + + def get_forecast(self, lead_days: int = 7) -> Optional[pd.DataFrame]: + """Fetch forecast for Hong Kong from selected model.""" + params = { + **self.params, + "daily": [ + "temperature_2m_max", + "temperature_2m_min", + "temperature_2m_mean", + "precipitation_sum", + "precipitation_probability_max", + "rain_sum", + "wind_speed_10m_max", + "wind_gusts_10m_max", + "wind_direction_10m_dominant", + "shortwave_radiation_sum", + "et0_fao_evapotranspiration", + "weather_code", + ], + "hourly": [ + "temperature_2m", + "relative_humidity_2m", + "dew_point_2m", + "apparent_temperature", + "precipitation_probability", + "precipitation", + "rain", + "cloud_cover", + "cloud_cover_low", + "cloud_cover_mid", + "cloud_cover_high", + "wind_speed_10m", + "wind_speed_100m", + "wind_gusts_10m", + "wind_direction_10m", + "wind_direction_100m", + "surface_pressure", + "visibility", + ], + "past_days": 0, + "forecast_days": lead_days, + } + + if OPENMETEO_API_KEY: + params["apikey"] = OPENMETEO_API_KEY + + try: + responses = self.client.weather_api(self.BASE_URL + "forecast", params=params) + return self._parse_response(responses[0]) + except Exception as e: + print(f"Open-Meteo API error: {e}") + return None + + def _parse_response(self, response) -> pd.DataFrame: + """Parse Open-Meteo response into a DataFrame.""" + hourly = response.Hourly() + daily = response.Daily() + + self._hourly_vars = [ + "temperature_2m", "relative_humidity_2m", "dew_point_2m", + "apparent_temperature", "precipitation_probability", "precipitation", + "rain", "cloud_cover", "cloud_cover_low", "cloud_cover_mid", + "cloud_cover_high", "wind_speed_10m", "wind_speed_100m", + "wind_gusts_10m", "wind_direction_10m", "wind_direction_100m", + "surface_pressure", "visibility", + ] + self._daily_vars = [ + "temperature_2m_max", "temperature_2m_min", "temperature_2m_mean", + "precipitation_sum", "precipitation_probability_max", "rain_sum", + "wind_speed_10m_max", "wind_gusts_10m_max", + "wind_direction_10m_dominant", "shortwave_radiation_sum", + "et0_fao_evapotranspiration", "weather_code", + ] + + hourly_data = { + "date": pd.date_range( + start=pd.Timestamp(hourly.Time(), unit="s", tz="UTC"), + end=pd.Timestamp(hourly.TimeEnd(), unit="s", tz="UTC"), + freq=pd.Timedelta(seconds=hourly.Interval()), + inclusive="left", + ) + } + for i in range(hourly.VariablesLength()): + var = hourly.Variables(i) + if i < len(self._hourly_vars): + hourly_data[self._hourly_vars[i]] = var.ValuesAsNumpy() + + hourly_df = pd.DataFrame(hourly_data).set_index("date") + hourly_df.index = hourly_df.index.tz_convert("Asia/Hong_Kong") + + daily_data = { + "date": pd.date_range( + start=pd.Timestamp(daily.Time(), unit="s", tz="UTC"), + end=pd.Timestamp(daily.TimeEnd(), unit="s", tz="UTC"), + freq=pd.Timedelta(seconds=daily.Interval()), + inclusive="left", + ) + } + for i in range(daily.VariablesLength()): + var = daily.Variables(i) + if i < len(self._daily_vars): + daily_data[self._daily_vars[i]] = var.ValuesAsNumpy() + + daily_df = pd.DataFrame(daily_data).set_index("date") + daily_df.index = daily_df.index.tz_convert("Asia/Hong_Kong") + + daily_df.attrs["model"] = self.model + daily_df.attrs["hourly"] = hourly_df + daily_df.attrs["fetch_time"] = datetime.now() + + return daily_df + + def get_current_conditions(self) -> dict: + """Get current weather conditions at HKO headquarters.""" + self._current_vars = [ + "temperature_2m", "relative_humidity_2m", "apparent_temperature", + "precipitation", "rain", "cloud_cover", "wind_speed_10m", + "wind_direction_10m", "wind_gusts_10m", "surface_pressure", + "weather_code", + ] + params = { + **self.params, + "current": self._current_vars, + "forecast_days": 1, + } + + try: + responses = self.client.weather_api(self.BASE_URL + "forecast", params=params) + current = responses[0].Current() + result = {"timestamp": datetime.now().isoformat()} + for i in range(current.VariablesLength()): + if i < len(self._current_vars): + result[self._current_vars[i]] = current.Variables(i).Value() + return { + "temperature": result.get("temperature_2m"), + "humidity": result.get("relative_humidity_2m"), + "apparent_temp": result.get("apparent_temperature"), + "precipitation": result.get("precipitation"), + "rain": result.get("rain"), + "cloud_cover": result.get("cloud_cover"), + "wind_speed": result.get("wind_speed_10m"), + "wind_direction": result.get("wind_direction_10m"), + "wind_gusts": result.get("wind_gusts_10m"), + "surface_pressure": result.get("surface_pressure"), + "weather_code": result.get("weather_code"), + "timestamp": datetime.now().isoformat(), + } + except Exception as e: + print(f"Open-Meteo current conditions error: {e}") + return {} + + def get_precipitation_probability(self, hours_ahead: int = 24) -> float: + """Get precipitation probability for next N hours.""" + df = self.get_forecast(lead_days=2) + if df is not None and hasattr(df, 'attrs') and 'hourly' in df.attrs: + hourly = df.attrs['hourly'] + now = pd.Timestamp.now(tz="Asia/Hong_Kong") + future = hourly[hourly.index <= now + pd.Timedelta(hours=hours_ahead)] + if 'precipitation_probability' in future.columns: + return float(future['precipitation_probability'].max()) + return 0.0 + + def get_scoring_window_summary(self, target_date: str) -> dict: + """ + Get a forecast summary for a specific scoring window (target date). + Used by the signal generator to create trading signals. + + Returns a dict with all relevant forecast variables for market comparison. + """ + df = self.get_forecast(lead_days=7) + if df is None: + return {} + + target = pd.Timestamp(target_date).tz_localize("Asia/Hong_Kong") + if target not in df.index: + return {} + + row = df.loc[target] + hourly = df.attrs.get("hourly", pd.DataFrame()) + + if not hourly.empty: + day_hourly = hourly[ + (hourly.index >= target) & + (hourly.index < target + pd.Timedelta(days=1)) + ] + else: + day_hourly = pd.DataFrame() + + summary = { + "date": target_date, + "model": self.model, + "fetch_time": df.attrs.get("fetch_time", datetime.now()).isoformat(), + "temperature_2m_max": float(row.get("temperature_2m_max", np.nan)), + "temperature_2m_min": float(row.get("temperature_2m_min", np.nan)), + "precipitation_sum": float(row.get("precipitation_sum", 0)), + "precipitation_probability_max": float(row.get("precipitation_probability_max", 0)), + "wind_speed_10m_max": float(row.get("wind_speed_10m_max", np.nan)), + "wind_gusts_10m_max": float(row.get("wind_gusts_10m_max", np.nan)), + } + + if not day_hourly.empty: + summary["temperature_2m_max_hourly"] = float(day_hourly["temperature_2m"].max()) if "temperature_2m" in day_hourly else np.nan + summary["precipitation_probability_max_hourly"] = float(day_hourly["precipitation_probability"].max()) if "precipitation_probability" in day_hourly else 0 + summary["precipitation_sum_hourly"] = float(day_hourly["precipitation"].sum()) if "precipitation" in day_hourly else 0 + + return summary