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.
This commit is contained in:
@@ -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
|
||||
+15
@@ -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
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
3.10.17
|
||||
@@ -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
|
||||
+112
@@ -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()
|
||||
@@ -0,0 +1,6 @@
|
||||
"""Polymarket integration for HK weather prediction markets."""
|
||||
|
||||
from .polymarket_client import PolymarketClient
|
||||
from .trader import Trader
|
||||
|
||||
__all__ = ["PolymarketClient", "Trader"]
|
||||
@@ -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",
|
||||
}
|
||||
@@ -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}")
|
||||
+246
@@ -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()
|
||||
@@ -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()
|
||||
@@ -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()
|
||||
@@ -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"
|
||||
@@ -0,0 +1,7 @@
|
||||
"""Trading strategy components."""
|
||||
|
||||
from .calibrator import ProbabilityCalibrator
|
||||
from .kelly import KellyCriterion
|
||||
from .signals import SignalGenerator
|
||||
|
||||
__all__ = ["ProbabilityCalibrator", "KellyCriterion", "SignalGenerator"]
|
||||
@@ -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),
|
||||
}
|
||||
@@ -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))
|
||||
@@ -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)
|
||||
@@ -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"]
|
||||
@@ -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()
|
||||
@@ -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
|
||||
@@ -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),
|
||||
}
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user