"""Correlation-aware portfolio Kelly sizing. When trading multiple weather markets simultaneously, outcomes are correlated: - Rain ↔ Temperature (rain suppresses temperature) - Wind ↔ Typhoon (wind is a prerequisite for T3/T8) - Rain tomorrow ↔ Rain day-after-tomorrow (persistence) Simple Kelly treats each market independently, over-betting when outcomes are correlated. This module implements simultaneous Kelly using a covariance matrix of historical outcomes. f* = Σ⁻¹ μ (vector form, where Σ is the covariance of returns, μ is edge vector) Reference: MacLean, Thorp, Ziemba (2011) "The Kelly Capital Growth Investment Criterion" """ import json from pathlib import Path from typing import Dict, List, Optional, Tuple import numpy as np from config import DATA_DIR COV_DIR = Path(DATA_DIR) / "covariance" COV_DIR.mkdir(parents=True, exist_ok=True) class PortfolioKelly: """ Simultaneous Kelly sizing for correlated prediction market bets. Uses historical outcome correlation matrix to adjust individual Kelly fractions, preventing over-betting on correlated markets. Usage: pk = PortfolioKelly() pk.update_correlation({"rain_yes": True, "temp_above_30": False}) sizes = pk.simultaneous_kelly(edges, sigmas, bankroll=1000) """ def __init__(self, fraction: float = 0.25, max_position_per_market: float = 500.0): self.fraction = fraction self.max_position_per_market = max_position_per_market # Outcome tracking self.outcome_history: Dict[str, List[int]] = {} self.target_names: List[str] = [] self.correlation_matrix: Optional[np.ndarray] = None self.n_observations: int = 0 self._load_correlation() def record_outcomes(self, outcomes: Dict[str, bool]): """ Record resolved market outcomes. Parameters ---------- outcomes : dict {target_name: bool} e.g., {"temp_gt_30c_24h": True, "rain_gt_0mm_24h": False} """ for name, outcome in outcomes.items(): if name not in self.outcome_history: self.outcome_history[name] = [] self.outcome_history[name].append(1 if outcome else 0) # Align all vectors to same length lengths = [len(v) for v in self.outcome_history.values()] if lengths: min_len = min(lengths) if min_len > self.n_observations: self.n_observations = min_len self._recompute_correlation() def _recompute_correlation(self): """Compute correlation matrix from outcome history.""" self.target_names = list(self.outcome_history.keys()) if len(self.target_names) < 2 or self.n_observations < 10: self.correlation_matrix = None return # Build outcome matrix n = min(self.n_observations, min(len(v) for v in self.outcome_history.values())) O = np.zeros((n, len(self.target_names))) for j, name in enumerate(self.target_names): O[:, j] = self.outcome_history[name][:n] # Pearson correlation of binary outcomes self.correlation_matrix = np.corrcoef(O.T) # Save self._save_correlation() def simultaneous_kelly( self, edges: Dict[str, float], sigmas: Optional[Dict[str, float]] = None, bankroll: float = 1000.0, ) -> Dict[str, float]: """ Compute simultaneous Kelly bet sizes. Parameters ---------- edges : dict {target_name: edge_in_decimal} e.g., {"temp_gt_30c": 0.15} sigmas : dict, optional {target_name: outcome_std} — if None, uses sqrt(p * (1-p)) bankroll : float Current bankroll Returns ------- dict {target_name: bet_size_in_usdc} """ targets = list(edges.keys()) n = len(targets) if n == 0: return {} # Edge vector mu = np.array([edges[t] for t in targets]) # Variance vector (binary outcome variance) if sigmas: sigma_diag = np.array([sigmas[t] for t in targets]) else: # Binary variance: p(1-p) for outcomes at price p sigma_diag = np.array([0.25] * n) # Max variance at p=0.5 # Build covariance matrix if self.correlation_matrix is not None and n >= 2: idxs = [] for t in targets: if t in self.target_names: idxs.append(self.target_names.index(t)) else: idxs.append(None) Sigma = np.zeros((n, n)) for i in range(n): for j in range(n): if i == j: Sigma[i, j] = sigma_diag[i] elif idxs[i] is not None and idxs[j] is not None: rho = self.correlation_matrix[idxs[i], idxs[j]] Sigma[i, j] = rho * np.sqrt(sigma_diag[i] * sigma_diag[j]) else: Sigma[i, j] = 0.0 # Unknown correlation → assume 0 else: # No correlation data → diagonal Sigma = np.diag(sigma_diag) # Regularize to ensure invertibility Sigma += np.eye(n) * 1e-6 try: # Simultaneous Kelly: f* = Σ⁻¹ μ f_star = np.linalg.solve(Sigma, mu) # Apply fractional Kelly f_fractional = f_star * self.fraction # Convert to USD sizes sizes = {} for i, name in enumerate(targets): size = f_fractional[i] * bankroll # Cap at per-market max, floor at 0 size = max(0, min(size, self.max_position_per_market)) sizes[name] = float(size) return sizes except np.linalg.LinAlgError: # Singular matrix fallback: independent Kelly sizes = {} for name in targets: size = (edges[name] * bankroll * self.fraction / max(sigma_diag[list(edges.keys()).index(name)], 0.01)) sizes[name] = float(max(0, min(size, self.max_position_per_market))) return sizes def get_correlation(self, target_a: str, target_b: str) -> float: """Get correlation between two target variables.""" if self.correlation_matrix is None: return 0.0 try: i = self.target_names.index(target_a) j = self.target_names.index(target_b) return float(self.correlation_matrix[i, j]) except (ValueError, IndexError): return 0.0 def correlation_summary(self) -> str: """Human-readable correlation summary.""" if self.correlation_matrix is None or len(self.target_names) < 2: return "No correlation data (need 10+ observations on 2+ targets)" lines = ["Correlation Matrix:", " " + " ".join(f"{n:>12s}" for n in self.target_names)] for i, name in enumerate(self.target_names): row = f" {name:<12s}" for j in range(len(self.target_names)): row += f"{self.correlation_matrix[i,j]:>12.3f}" lines.append(row) return "\n".join(lines) def _save_correlation(self): """Persist correlation data.""" path = COV_DIR / "correlation.json" data = { "target_names": self.target_names, "n_observations": self.n_observations, "correlation_matrix": self.correlation_matrix.tolist() if self.correlation_matrix is not None else None, "outcome_history": self.outcome_history, } with open(path, "w") as f: json.dump(data, f) def _load_correlation(self): """Load persisted correlation data.""" path = COV_DIR / "correlation.json" if not path.exists(): return try: with open(path) as f: data = json.load(f) self.target_names = data.get("target_names", []) self.n_observations = data.get("n_observations", 0) self.outcome_history = data.get("outcome_history", {}) corr = data.get("correlation_matrix") if corr is not None: self.correlation_matrix = np.array(corr) except Exception: pass def independent_kelly( self, edges: Dict[str, float], bankroll: float = 1000.0, ) -> Dict[str, float]: """Independent (simple) Kelly sizing — for comparison.""" sizes = {} for name, edge in edges.items(): # Simple Kelly: f* = edge / sigma^2 for binary outcomes sigma_sq = 0.25 f_star = edge / max(sigma_sq, 0.01) size = f_star * self.fraction * bankroll sizes[name] = float(max(0, min(size, self.max_position_per_market))) return sizes def compare(self, edges: Dict[str, float], bankroll: float = 1000.0) -> Dict: """Compare independent vs simultaneous Kelly allocations.""" ind = self.independent_kelly(edges, bankroll) sim = self.simultaneous_kelly(edges, bankroll=bankroll) comparison = {} for name in edges: comparison[name] = { "independent": ind.get(name, 0), "simultaneous": sim.get(name, 0), "reduction": ind.get(name, 0) - sim.get(name, 0), } return comparison