Files
ramseshk c93af97059 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.
2026-08-10 12:48:05 +08:00

247 lines
9.1 KiB
Python

#!/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()