0446443d36
New strategies: - Cross-Sectional Momentum: long top-N, short bottom-N across HL universe - Spot-Perp Basis Arbitrage: delta-neutral spot vs perp price gap trading - Regime-Switching Ensemble: dynamically allocates strategies by market regime - Portfolio Construction: risk parity, vol targeting, correlation penalty Infrastructure: - DuckDBDataProvider: real tick/candle data for backtests (replaces synthetic) - Walk-Forward Validation: systematic IS/OOS across all 12 strategies - 3 Jupyter research notebooks (EDA, strategy research, portfolio) Pipeline integration: - deploy.py registry, sweep_runner, vbt_runner all updated - fee_tiers support for new strategies - All modules syntax-validated and import-tested
421 lines
16 KiB
Python
421 lines
16 KiB
Python
"""
|
|
Portfolio construction layer — risk allocation across strategies.
|
|
|
|
Turns N independent strategy signals into a single meta-portfolio using:
|
|
1. Risk parity — allocates capital inversely proportional to strategy vol
|
|
2. Volatility targeting — scales total portfolio to target annualized vol
|
|
3. Correlation-based sizing — reduces allocation to redundant strategies
|
|
4. Maximum drawdown stops — kill switch per strategy and portfolio-level
|
|
5. Regime-weighted allocation — adjusts weights based on market regime
|
|
|
|
Integration point: sits between strategy signals and execution.
|
|
Consumes signal strength values from each strategy, produces position sizes.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
from collections import deque
|
|
from dataclasses import dataclass, field
|
|
|
|
import numpy as np
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
@dataclass
|
|
class StrategyAllocation:
|
|
"""Position and PnL state for one strategy within the portfolio."""
|
|
name: str
|
|
weight: float = 0.0
|
|
position: float = 0.0 # Current signed position (units)
|
|
entry_price: float = 0.0 # Average entry price
|
|
pnl: float = 0.0 # Realized PnL
|
|
unrealized: float = 0.0 # Mark-to-market PnL
|
|
trades: int = 0 # Trade count
|
|
wins: int = 0 # Winning trades
|
|
fee_paid: float = 0.0
|
|
equity: float = 0.0 # Current allocation value
|
|
initial_equity: float = 0.0 # Starting allocation
|
|
signal_strength: deque = field(default_factory=lambda: deque(maxlen=100))
|
|
returns: deque = field(default_factory=lambda: deque(maxlen=500))
|
|
vol_20d: float = 0.0 # Rolling 20-period volatility
|
|
var_95: float = 0.0 # Value at Risk (95%)
|
|
active: bool = True # Kill-switch: False = disabled
|
|
max_drawdown: float = 0.0 # Peak-to-trough drawdown
|
|
peak_equity: float = 0.0 # All-time high equity
|
|
regime_scores: dict = field(default_factory=dict)
|
|
|
|
|
|
@dataclass
|
|
class PortfolioMetrics:
|
|
"""Aggregate portfolio metrics."""
|
|
total_equity: float = 0.0
|
|
total_exposure: float = 0.0
|
|
gross_exposure: float = 0.0
|
|
net_exposure: float = 0.0
|
|
total_pnl: float = 0.0
|
|
total_pnl_pct: float = 0.0
|
|
sharpe: float = 0.0
|
|
sortino: float = 0.0
|
|
vol_20d: float = 0.0
|
|
var_95: float = 0.0
|
|
cvar_95: float = 0.0
|
|
max_drawdown_pct: float = 0.0
|
|
daily_drawdown: float = 0.0
|
|
trades_today: int = 0
|
|
win_rate: float = 0.0
|
|
correlation_matrix: dict = field(default_factory=dict)
|
|
regime: str = "NORMAL"
|
|
|
|
|
|
class PortfolioConstructor:
|
|
"""Risk-managed portfolio of strategies.
|
|
|
|
Responsibilities:
|
|
- Compute optimal capital allocation per strategy
|
|
- Apply volatility targeting at portfolio level
|
|
- Reduce allocations to correlated strategies
|
|
- Enforce per-strategy and portfolio-level drawdown stops
|
|
- Produce final position sizes for each strategy
|
|
|
|
Usage:
|
|
pf = PortfolioConstructor(capital=100000, vol_target=0.20, max_correlation=0.70)
|
|
pf.update_returns("pairs", [0.001, -0.002, 0.003])
|
|
pf.update_signal("pairs", strength=0.8, direction="BUY")
|
|
...
|
|
sizes = pf.get_positions(current_prices)
|
|
"""
|
|
|
|
def __init__(
|
|
self,
|
|
capital: float = 100_000.0,
|
|
vol_target: float = 0.20, # Annualized vol target
|
|
max_correlation: float = 0.70, # Max corr before reducing allocation
|
|
min_allocation: float = 0.02, # Min allocation fraction
|
|
max_allocation: float = 0.25, # Max allocation fraction per strategy
|
|
max_drawdown_stop: float = 0.15, # Kill strategy after 15% DD
|
|
portfolio_mdd_stop: float = 0.10, # Stop entire portfolio at 10% DD
|
|
n_lookback: int = 200, # Days for risk estimation
|
|
regime_weights: dict | None = None, # Per-regime strategy weights
|
|
):
|
|
self.capital = capital
|
|
self.vol_target = vol_target
|
|
self.max_correlation = max_correlation
|
|
self.min_allocation = min_allocation
|
|
self.max_allocation = max_allocation
|
|
self.max_drawdown_stop = max_drawdown_stop
|
|
self.portfolio_mdd_stop = portfolio_mdd_stop
|
|
self.n_lookback = n_lookback
|
|
|
|
self.strategies: dict[str, StrategyAllocation] = {}
|
|
self.portfolio_returns: deque = deque(maxlen=n_lookback)
|
|
self.portfolio_equity_history: deque = deque(maxlen=n_lookback)
|
|
self.peak_equity: float = capital
|
|
self.current_regime: str = "NORMAL"
|
|
self.regime_weights = regime_weights or {}
|
|
|
|
self._portfolio_stopped: bool = False
|
|
|
|
def register_strategy(self, name: str, allocation: float = 0.0):
|
|
"""Register a strategy in the portfolio."""
|
|
if name not in self.strategies:
|
|
alloc = allocation if allocation > 0 else self.capital * self.min_allocation
|
|
self.strategies[name] = StrategyAllocation(
|
|
name=name,
|
|
initial_equity=alloc,
|
|
equity=alloc,
|
|
weight=1.0 / max(len(self.strategies) + 1, 1),
|
|
)
|
|
|
|
def update_returns(self, name: str, returns: list[float]):
|
|
"""Feed per-bar returns for a strategy."""
|
|
if name not in self.strategies:
|
|
self.register_strategy(name)
|
|
st = self.strategies[name]
|
|
for r in returns:
|
|
st.returns.append(r)
|
|
|
|
def update_signal(self, name: str, strength: float, direction: str,
|
|
price: float = 0.0, regime: str = "NORMAL"):
|
|
"""Record a strategy signal and its strength."""
|
|
if name not in self.strategies:
|
|
self.register_strategy(name)
|
|
st = self.strategies[name]
|
|
st.signal_strength.append(strength)
|
|
if regime not in st.regime_scores:
|
|
st.regime_scores[regime] = []
|
|
st.regime_scores[regime].append(strength if direction == "BUY" else -strength)
|
|
|
|
def compute_allocations(self, current_prices: dict[str, float]) -> dict[str, float]:
|
|
"""Compute optimal capital allocation per strategy.
|
|
|
|
Returns dict of strategy_name → dollar_amount to allocate.
|
|
"""
|
|
total_alloc = 0.0
|
|
weights = {}
|
|
vols = {}
|
|
n_active = sum(1 for s in self.strategies.values() if s.active)
|
|
|
|
# Step 1: compute individual strategy vols
|
|
for name, st in self.strategies.items():
|
|
if not st.active or len(st.returns) < 20:
|
|
weights[name] = 0.0
|
|
continue
|
|
returns = list(st.returns)[-min(len(st.returns), self.n_lookback):]
|
|
vol = float(np.std(returns)) if len(returns) > 1 else 0.0
|
|
annual_vol = vol * np.sqrt(365 * 24) if vol > 0 else 0.10
|
|
st.vol_20d = annual_vol
|
|
vols[name] = annual_vol
|
|
|
|
# VaR 95%
|
|
if len(returns) >= 50:
|
|
st.var_95 = float(np.percentile(returns, 5))
|
|
|
|
weights[name] = 1.0 / max(annual_vol, 0.01)
|
|
|
|
if not vols:
|
|
return {name: 0.0 for name in self.strategies}
|
|
|
|
# Step 2: adjust weights for correlation — reduce allocation to highly correlated strategies
|
|
adjusted_weights = self._adjust_for_correlation(weights)
|
|
|
|
# Step 3: normalize to sum to 1 (risk-parity)
|
|
total = sum(adjusted_weights.values())
|
|
if total > 0:
|
|
for name in adjusted_weights:
|
|
adjusted_weights[name] = max(
|
|
self.min_allocation,
|
|
min(self.max_allocation, adjusted_weights[name] / total)
|
|
)
|
|
|
|
# Step 4: regime override — if regime weights are specified, blend with risk-parity
|
|
if self.current_regime in self.regime_weights:
|
|
rw = self.regime_weights[self.current_regime]
|
|
for name, wt in rw.items():
|
|
if name in adjusted_weights:
|
|
adjusted_weights[name] = adjusted_weights.get(name, 0) * 0.5 + wt * 0.5
|
|
|
|
# Step 5: vol target scaling
|
|
if self.vol_target > 0 and self.portfolio_returns:
|
|
pf_returns = list(self.portfolio_returns)[-100:]
|
|
if len(pf_returns) > 10:
|
|
pf_vol = float(np.std(pf_returns)) * np.sqrt(365 * 24)
|
|
scale = self.vol_target / max(pf_vol, 0.01)
|
|
scale = min(scale, 2.0) # Max 2x leverage
|
|
for name in adjusted_weights:
|
|
adjusted_weights[name] *= scale
|
|
|
|
# Step 6: convert weights to dollar allocations
|
|
allocations = {}
|
|
working_capital = self.capital * 0.70 # 70% of capital deployed, 30% reserve
|
|
total_weight = sum(adjusted_weights.values())
|
|
for name, wt in adjusted_weights.items():
|
|
if total_weight > 0:
|
|
allocations[name] = working_capital * (wt / total_weight)
|
|
else:
|
|
allocations[name] = 0.0
|
|
|
|
for name, st in self.strategies.items():
|
|
st.weight = adjusted_weights.get(name, 0.0)
|
|
st.equity = allocations.get(name, 0.0)
|
|
|
|
return allocations
|
|
|
|
def _adjust_for_correlation(self, raw_weights: dict[str, float]) -> dict[str, float]:
|
|
"""Reduce weights of correlated strategies to avoid concentration."""
|
|
adjusted = dict(raw_weights)
|
|
|
|
strategy_names = [n for n in raw_weights if raw_weights[n] > 0]
|
|
if len(strategy_names) < 2:
|
|
return adjusted
|
|
|
|
# Build correlation matrix from returns
|
|
corr_matrix = {}
|
|
for i, n1 in enumerate(strategy_names):
|
|
for n2 in strategy_names[i + 1:]:
|
|
r1 = list(self.strategies[n1].returns)[-200:]
|
|
r2 = list(self.strategies[n2].returns)[-200:]
|
|
min_len = min(len(r1), len(r2))
|
|
if min_len < 20:
|
|
corr = 0.0
|
|
else:
|
|
corr = float(np.corrcoef(r1[-min_len:], r2[-min_len:])[0, 1])
|
|
corr_matrix[f"{n1}|{n2}"] = round(corr, 3)
|
|
|
|
# Penalize correlated pairs
|
|
for key, corr in corr_matrix.items():
|
|
if abs(corr) > self.max_correlation and not np.isnan(corr):
|
|
n1, n2 = key.split("|")
|
|
penalty = 1.0 - (abs(corr) - self.max_correlation)
|
|
adjusted[n1] = adjusted.get(n1, 0) * penalty
|
|
adjusted[n2] = adjusted.get(n2, 0) * penalty
|
|
|
|
return adjusted
|
|
|
|
def get_positions(self, current_prices: dict[str, float],
|
|
signals: dict[str, dict] | None = None) -> dict[str, dict]:
|
|
"""Compute final position sizes for each strategy.
|
|
|
|
Args:
|
|
current_prices: coin → current mark price
|
|
signals: strategy_name → {"direction": "BUY"/"SELL", "strength": 0.0}
|
|
|
|
Returns:
|
|
strategy_name → {"coin": str, "side": "BUY"/"SELL", "size": float, "price": float}
|
|
"""
|
|
allocations = self.compute_allocations(current_prices)
|
|
positions = {}
|
|
|
|
for name, alloc in allocations.items():
|
|
if alloc <= 0 or name not in self.strategies:
|
|
continue
|
|
st = self.strategies[name]
|
|
if not st.active:
|
|
continue
|
|
|
|
# Determine coin from strategy name
|
|
coin = self._strategy_coin(name)
|
|
px = current_prices.get(coin, 0)
|
|
if px <= 0:
|
|
continue
|
|
|
|
# Signal-based direction override
|
|
sig = signals.get(name) if signals else None
|
|
if sig:
|
|
direction = sig.get("direction", "NEUTRAL")
|
|
strength = sig.get("strength", 0.0)
|
|
else:
|
|
# Default: use recent signal history
|
|
recent = list(st.signal_strength)[-20:]
|
|
avg_signal = float(np.mean(recent)) if recent else 0.0
|
|
direction = "BUY" if avg_signal > 0 else "SELL"
|
|
strength = abs(avg_signal)
|
|
|
|
# Position size: allocation / price, scaled by signal strength
|
|
base_size = alloc / px
|
|
size = base_size * min(strength, 1.5)
|
|
size = max(size, base_size * 0.25) # Minimum 25% of base size
|
|
|
|
positions[name] = {
|
|
"coin": coin,
|
|
"side": direction if strength > 0.1 else "NEUTRAL",
|
|
"size": round(size, 6),
|
|
"price": px,
|
|
"allocation": round(alloc, 2),
|
|
"weight": round(st.weight, 3),
|
|
"vol_20d": round(st.vol_20d, 3),
|
|
}
|
|
|
|
return positions
|
|
|
|
def _strategy_coin(self, name: str) -> str:
|
|
"""Map strategy to primary trading coin."""
|
|
coin_map = {
|
|
"pairs": "ETH",
|
|
"hurst_vpin": "BTC",
|
|
"as_mm": "BTC",
|
|
"obi": "BTC",
|
|
"grid_mm": "BTC",
|
|
"composite_mm": "BTC",
|
|
"iceberg": "BTC",
|
|
"funding_arb": "BTC",
|
|
"momentum": "ETH",
|
|
"mean_rev": "ETH",
|
|
"cross_sectional": "BTC",
|
|
}
|
|
return coin_map.get(name.lower(), "BTC")
|
|
|
|
def check_drawdown_stops(self) -> dict[str, bool]:
|
|
"""Check and enforce drawdown stops.
|
|
|
|
Returns dict of strategy_name → stopped (True if kill switch triggered).
|
|
"""
|
|
stops = {}
|
|
|
|
for name, st in self.strategies.items():
|
|
if not st.active:
|
|
continue
|
|
|
|
if st.equity > st.peak_equity:
|
|
st.peak_equity = st.equity
|
|
|
|
if st.peak_equity > 0:
|
|
dd = 1.0 - st.equity / st.peak_equity
|
|
if dd > self.max_drawdown_stop:
|
|
st.active = False
|
|
stops[name] = True
|
|
logger.warning("KILL SWITCH: %s DD=%.1f%% > %.1f%% limit",
|
|
name, dd * 100, self.max_drawdown_stop * 100)
|
|
else:
|
|
stops[name] = False
|
|
|
|
return stops
|
|
|
|
def update_portfolio_value(self, strategy_pnls: dict[str, float]):
|
|
"""Update portfolio equity after a round of PnL."""
|
|
total_pnl = sum(strategy_pnls.values())
|
|
new_equity = self.capital + total_pnl
|
|
|
|
if len(self.portfolio_equity_history) > 0:
|
|
prev = self.portfolio_equity_history[-1]
|
|
if prev > 0:
|
|
ret = (new_equity - prev) / prev
|
|
self.portfolio_returns.append(ret)
|
|
|
|
self.portfolio_equity_history.append(new_equity)
|
|
|
|
if new_equity > self.peak_equity:
|
|
self.peak_equity = new_equity
|
|
|
|
# Portfolio-level drawdown stop
|
|
if self.peak_equity > 0:
|
|
pf_dd = 1.0 - new_equity / self.peak_equity
|
|
if pf_dd > self.portfolio_mdd_stop and not self._portfolio_stopped:
|
|
self._portfolio_stopped = True
|
|
logger.warning("PORTFOLIO KILL SWITCH: DD=%.1f%% > %.1f%%",
|
|
pf_dd * 100, self.portfolio_mdd_stop * 100)
|
|
|
|
def summary(self) -> PortfolioMetrics:
|
|
"""Generate portfolio metrics report."""
|
|
pf = PortfolioMetrics()
|
|
pf.regime = self.current_regime
|
|
|
|
equity_vals = list(self.portfolio_equity_history)
|
|
if equity_vals:
|
|
pf.total_equity = round(equity_vals[-1], 2)
|
|
returns = list(self.portfolio_returns)
|
|
if len(returns) > 10:
|
|
pf.vol_20d = round(float(np.std(returns[-20:])) * np.sqrt(365 * 24), 3)
|
|
pf.sharpe = round(float(np.mean(returns) / max(np.std(returns), 1e-10)) * np.sqrt(365 * 24), 3)
|
|
down_returns = [r for r in returns if r < 0]
|
|
if down_returns:
|
|
pf.sortino = round(float(np.mean(returns) / max(np.std(down_returns), 1e-10)) * np.sqrt(365 * 24), 3)
|
|
|
|
# Max drawdown
|
|
peak = equity_vals[0]
|
|
pf.max_drawdown_pct = 0.0
|
|
for v in equity_vals:
|
|
if v > peak:
|
|
peak = v
|
|
dd = (peak - v) / peak if peak > 0 else 0
|
|
if dd > pf.max_drawdown_pct:
|
|
pf.max_drawdown_pct = dd
|
|
|
|
pf.total_pnl = sum(s.pnl for s in self.strategies.values())
|
|
pf.total_pnl_pct = pf.total_pnl / self.capital if self.capital > 0 else 0
|
|
|
|
pf.trades_today = sum(s.trades for s in self.strategies.values())
|
|
total_wins = sum(s.wins for s in self.strategies.values())
|
|
pf.win_rate = total_wins / max(pf.trades_today, 1)
|
|
|
|
long_exp = sum(s.equity for n, s in self.strategies.items() if s.position > 0)
|
|
short_exp = sum(abs(s.position) for n, s in self.strategies.items() if s.position < 0)
|
|
pf.gross_exposure = long_exp + short_exp
|
|
pf.net_exposure = long_exp - short_exp
|
|
pf.total_exposure = pf.gross_exposure
|
|
|
|
return pf
|
|
|
|
def is_stopped(self) -> bool:
|
|
return self._portfolio_stopped
|