-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathensemble_integration.py
More file actions
324 lines (276 loc) · 13.4 KB
/
Copy pathensemble_integration.py
File metadata and controls
324 lines (276 loc) · 13.4 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
"""
Ensemble Integration Module
Manages adaptive strategy weighting and regime-aware ensemble optimization
"""
import asyncio
import logging
from typing import Dict, Any, List, Optional
import pandas as pd
import numpy as np
from mathematical_ensemble_optimiser import (
EnsembleOptimizer,
OptimizationConfig,
OptimizationObjective,
MarketRegime
)
try:
from crypto_ensemble_optimizer import (
CryptoEnsembleOptimizer,
CryptoOptimizationConfig
)
CRYPTO_OPTIMIZER_AVAILABLE = True
except ImportError:
logger.warning("Crypto optimizer not available")
CRYPTO_OPTIMIZER_AVAILABLE = False
logger = logging.getLogger(__name__)
class EnhancedEnsembleManager:
"""
Manages the full ensemble of strategies with adaptive weighting and regime detection.
Integrates with MirrorCore-X SyncBus architecture.
"""
def __init__(self, strategy_trainer, sync_bus, risk_profile: str = 'moderate'):
self.strategy_trainer = strategy_trainer
self.sync_bus = sync_bus
self.risk_profile = risk_profile
self.adaptive_enabled = True
# Import adaptive optimizer
try:
from advanced_strategies import AdaptiveEnsembleOptimizer
strategies = list(strategy_trainer.strategies.values())
self.ensemble_optimizer = AdaptiveEnsembleOptimizer(strategies, risk_profile)
self.optimizer_available = True
except ImportError:
logger.warning("Adaptive ensemble optimizer not available")
self.ensemble_optimizer = None
self.optimizer_available = False
# Initialize crypto optimizer if available
self.crypto_optimizer = None
if CRYPTO_OPTIMIZER_AVAILABLE:
try:
crypto_config = CryptoOptimizationConfig(
risk_profile=risk_profile,
objective=OptimizationObjective.SHARPE_RATIO,
max_position_size=0.1,
turnover_penalty=0.01
)
self.crypto_optimizer = CryptoEnsembleOptimizer(
strategies=strategies,
config=crypto_config
)
logger.info("Crypto ensemble optimizer initialized successfully.")
except Exception as e:
logger.error(f"Failed to initialize crypto optimizer: {e}")
self.crypto_optimizer = None
self.current_weights = {}
self.current_regime = 'normal'
self.performance_tracking = {}
async def update(self, data: Dict[str, Any]) -> Dict[str, Any]:
"""Update ensemble weights based on current market regime and performance"""
try:
scanner_data = data.get('scanner_data', [])
if not scanner_data:
return {}
df = pd.DataFrame(scanner_data)
# Detect current regime
if self.ensemble_optimizer:
self.current_regime = self.ensemble_optimizer.detect_regime(df)
# Calculate strategy performance
recent_performance = self._calculate_recent_performance()
# Use mathematical optimization for weights
if hasattr(self.strategy_trainer, 'optimize_ensemble_weights'):
math_weights = self.strategy_trainer.optimize_ensemble_weights(market_data=df)
# Blend mathematical weights with adaptive weights
if self.ensemble_optimizer and self.adaptive_enabled:
adaptive_weights = self.ensemble_optimizer.calculate_weights(
regime=self.current_regime,
recent_performance=recent_performance
)
# 70% mathematical optimization, 30% adaptive heuristic
self.current_weights = {
name: 0.7 * math_weights.get(name, 0) + 0.3 * adaptive_weights.get(name, 0)
for name in set(list(math_weights.keys()) + list(adaptive_weights.keys()))
}
else:
self.current_weights = math_weights
else:
# Fallback to adaptive weights only
if self.ensemble_optimizer and self.adaptive_enabled:
self.current_weights = self.ensemble_optimizer.calculate_weights(
regime=self.current_regime,
recent_performance=recent_performance
)
# If crypto optimizer is available, use it to refine weights
if self.crypto_optimizer:
try:
# Assuming market_data passed to crypto optimizer includes relevant crypto features
# We might need to adjust the data passed based on crypto_optimizer's requirements
crypto_weights = self.crypto_optimizer.optimize(
market_data=df,
current_weights=self.current_weights,
current_regime=self.current_regime,
recent_performance=recent_performance
)
# Blend traditional weights with crypto-optimized weights
# Example: 60% traditional, 40% crypto
self.current_weights = {
name: 0.6 * self.current_weights.get(name, 0) + 0.4 * crypto_weights.get(name, 0)
for name in set(list(self.current_weights.keys()) + list(crypto_weights.keys()))
}
logger.info("Ensemble weights refined by crypto optimizer.")
except Exception as e:
logger.error(f"Crypto optimizer failed to refine weights: {e}")
# Update global state
await self.sync_bus.update_state('ensemble_weights', self.current_weights)
await self.sync_bus.update_state('market_regime', self.current_regime)
# Get optimization report if available
opt_report = {}
if hasattr(self.strategy_trainer, 'get_optimization_report'):
opt_report = self.strategy_trainer.get_optimization_report()
logger.info(f"Ensemble updated: regime={self.current_regime}, weights={len(self.current_weights)}")
return {
'regime': self.current_regime,
'weights': self.current_weights,
'performance': recent_performance,
'optimization_report': opt_report
}
except Exception as e:
logger.error(f"Ensemble update failed: {e}")
return {}
def _calculate_recent_performance(self) -> Dict[str, Dict[str, float]]:
"""Calculate recent performance metrics for each strategy"""
performance = {}
for strategy_name, perf_data in self.strategy_trainer.performance_tracker.items():
if not perf_data:
continue
recent = perf_data[-20:] # Last 20 trades
if recent:
wins = len([p for p in recent if p > 0])
avg_return = np.mean(recent)
volatility = np.std(recent) if len(recent) > 1 else 0.5
sharpe = (avg_return / volatility) if volatility > 0 else 0
performance[strategy_name] = {
'win_rate': wins / len(recent),
'sharpe': sharpe,
'volatility': volatility,
'avg_return': avg_return
}
return performance
async def generate_ensemble_signal(self, df: pd.DataFrame) -> Dict[str, Any]:
"""Generate consensus signal from all strategies"""
if not self.ensemble_optimizer and not self.crypto_optimizer:
return {'direction': 'HOLD', 'strength': 0, 'confidence': 0}
# Get signals from all strategies
signals = {}
for strategy_name, agent in self.strategy_trainer.strategies.items():
try:
signal_value = agent.evaluate(df) if hasattr(agent, 'evaluate') else 0
signals[strategy_name] = float(signal_value)
except Exception as e:
logger.error(f"Strategy {strategy_name} evaluation failed: {e}")
signals[strategy_name] = 0.0
# Aggregate with current weights
if self.crypto_optimizer:
# Use crypto optimizer for aggregation if available, as it might handle crypto-specific aspects
consensus = self.crypto_optimizer.aggregate_signals(signals, self.current_weights)
elif self.ensemble_optimizer:
consensus = self.ensemble_optimizer.aggregate_signals(signals, self.current_weights)
else:
# Fallback if no optimizer is available (should not happen with current setup)
consensus = {'direction': 'HOLD', 'strength': 0, 'confidence': 0}
return consensus
def get_status(self) -> Dict[str, Any]:
"""Get ensemble status"""
return {
'optimizer_available': self.optimizer_available,
'crypto_optimizer_available': CRYPTO_OPTIMIZER_AVAILABLE,
'adaptive_enabled': self.adaptive_enabled,
'current_regime': self.current_regime,
'risk_profile': self.risk_profile,
'active_strategies': len(self.current_weights),
'weights': self.current_weights
}
async def create_enhanced_ensemble(strategy_trainer, sync_bus, risk_profile: str = 'moderate'):
"""Factory function to create enhanced ensemble manager"""
ensemble = EnhancedEnsembleManager(strategy_trainer, sync_bus, risk_profile)
# Attach to SyncBus
sync_bus.attach('ensemble_manager', ensemble)
logger.info(f"Enhanced ensemble manager created with {risk_profile} risk profile")
return ensemble
def get_optimal_weights(
lambda_risk: float = 100.0,
eta_turnover: float = 0.05,
max_weight: float = 0.25,
regime: str = "trending",
use_shrinkage: bool = True,
use_resampling: bool = False
):
"""Get optimal weights for current strategy ensemble"""
try:
# Get strategy list with categories
strategies = [
{'name': 'UT_BOT', 'category': 'trend', 'sharpe': 1.2},
{'name': 'GRADIENT_TREND', 'category': 'trend', 'sharpe': 1.15},
{'name': 'VOLUME_SR', 'category': 'support', 'sharpe': 1.3},
{'name': 'MEAN_REVERSION', 'category': 'reversion', 'sharpe': 1.4},
{'name': 'MOMENTUM_BREAKOUT', 'category': 'breakout', 'sharpe': 1.1},
{'name': 'VOLATILITY_REGIME', 'category': 'adaptive', 'sharpe': 1.25},
{'name': 'PAIRS_TRADING', 'category': 'arbitrage', 'sharpe': 1.35},
{'name': 'ANOMALY_DETECTION', 'category': 'ml', 'sharpe': 1.2},
{'name': 'SENTIMENT_MOMENTUM', 'category': 'hybrid', 'sharpe': 1.3},
{'name': 'REGIME_CHANGE', 'category': 'adaptive', 'sharpe': 1.15},
{'name': 'BAYESIAN_BELIEF', 'category': 'probabilistic', 'sharpe': 1.65},
{'name': 'LIQUIDITY_FLOW', 'category': 'hybrid', 'sharpe': 1.45},
{'name': 'MARKET_ENTROPY', 'category': 'information', 'sharpe': 1.15},
]
# Regime-based weight adjustments
regime_multipliers = {
'trending': {'UT_BOT': 2.5, 'GRADIENT_TREND': 2.2, 'MEAN_REVERSION': 0.3},
'ranging': {'MEAN_REVERSION': 3.0, 'PAIRS_TRADING': 2.8, 'UT_BOT': 0.5},
'volatile': {'VOLATILITY_REGIME': 3.0, 'ANOMALY_DETECTION': 2.5, 'MARKET_ENTROPY': 2.6},
'mixed': {}
}
multipliers = regime_multipliers.get(regime, {})
# Calculate base weights with regime adjustment
weights = []
for s in strategies:
mult = multipliers.get(s['name'], 1.0)
base_weight = (s['sharpe'] / 1.0) * mult
weights.append(min(base_weight, max_weight * len(strategies)))
# Normalize
total = sum(weights)
weights = [w / total for w in weights]
# Calculate portfolio metrics
expected_return = sum(w * s['sharpe'] * 0.001 for w, s in zip(weights, strategies))
expected_vol = 0.012 # Simplified
sharpe = expected_return / expected_vol if expected_vol > 0 else 0
return {
'weights': [
{'name': s['name'], 'weight': w, 'category': s['category']}
for s, w in zip(strategies, weights)
],
'expected_return': expected_return,
'expected_volatility': expected_vol,
'sharpe_ratio': sharpe,
'turnover': 0.05
}
except Exception as e:
logger.error(f"get_optimal_weights failed: {e}")
return {
'weights': [],
'expected_return': 0.0,
'expected_volatility': 0.0,
'sharpe_ratio': 0.0,
'turnover': 0.0
}
# Placeholder for demo function, assuming it exists elsewhere or is not needed for this specific change.
# If it's essential, it should be provided or defined.
def demo_ensemble_integration():
logger.info("Running demo ensemble integration (placeholder)")
pass # Placeholder for actual demo logic
if __name__ == "__main__":
logging.basicConfig(level=logging.INFO)
# Call the demo function if it's intended to be run directly
# For this example, we'll assume demo_ensemble_integration is defined elsewhere or not critical for the provided snippet.
# If it were critical, its definition would need to be included.
# demo_ensemble_integration()
logger.info("Ensemble integration module loaded.")