Move 17 library crates into crates/, CLI binary into bin/fxt, consolidate 10 test crates into testing/, split config crate from deployment config files. Root directory reduced from 38+ to ~17 directories. All Cargo.toml paths and build.rs proto refs updated. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
997 lines
38 KiB
Rust
997 lines
38 KiB
Rust
//! End-to-End Trading Workflows for E2E Testing
|
|
//!
|
|
//! Comprehensive trading workflow tests that validate the complete system:
|
|
//! - Order lifecycle from submission to execution
|
|
//! - Risk management and validation
|
|
//! - Market data processing and ML integration
|
|
//! - Backtesting workflows
|
|
//! - Error handling and recovery scenarios
|
|
|
|
#![allow(deprecated)]
|
|
|
|
use anyhow::{Context, Result};
|
|
use std::collections::HashMap;
|
|
use std::sync::Arc;
|
|
use std::time::{Duration, Instant};
|
|
use tokio::sync::RwLock;
|
|
use tokio::time::{sleep, timeout};
|
|
use tokio_stream::StreamExt;
|
|
use tracing::{debug, info, warn};
|
|
|
|
use crate::clients::TliClient;
|
|
use crate::database::TestDatabase;
|
|
use crate::ml_pipeline::MLTestPipeline;
|
|
use crate::utils::TestDataGenerator;
|
|
|
|
// Trading service proto types (for TradingWorkflow)
|
|
use crate::proto::trading::{
|
|
CancelOrderRequest, GetOrderStatusRequest, GetPortfolioSummaryRequest, MarketDataType,
|
|
OrderSide, OrderStatus, OrderType, StreamMarketDataRequest, StreamOrdersRequest,
|
|
SubmitOrderRequest,
|
|
};
|
|
|
|
// Backtesting service proto types (for BacktestingWorkflow)
|
|
use crate::proto::backtesting::{
|
|
GetBacktestResultsRequest, GetBacktestStatusRequest, ListBacktestsRequest,
|
|
StartBacktestRequest, StopBacktestRequest, SubscribeBacktestProgressRequest,
|
|
};
|
|
|
|
// Note: monitoring and risk proto modules are not implemented yet
|
|
// Future enhancement: Add monitoring and risk service proto definitions when needed
|
|
|
|
// Import missing types
|
|
use crate::ml_pipeline::PredictionType;
|
|
|
|
/// Trading workflow test result
|
|
#[derive(Debug, Clone)]
|
|
pub struct WorkflowTestResult {
|
|
pub workflow_name: String,
|
|
pub success: bool,
|
|
pub duration: Duration,
|
|
pub steps_completed: usize,
|
|
pub total_steps: usize,
|
|
pub error_message: Option<String>,
|
|
pub metrics: HashMap<String, f64>,
|
|
pub order_ids: Vec<String>,
|
|
pub trades_executed: usize,
|
|
}
|
|
|
|
impl WorkflowTestResult {
|
|
pub fn success(name: String, duration: Duration, steps: usize) -> Self {
|
|
Self {
|
|
workflow_name: name,
|
|
success: true,
|
|
duration,
|
|
steps_completed: steps,
|
|
total_steps: steps,
|
|
error_message: None,
|
|
metrics: HashMap::new(),
|
|
order_ids: Vec::new(),
|
|
trades_executed: 0,
|
|
}
|
|
}
|
|
|
|
pub fn failure(
|
|
name: String,
|
|
duration: Duration,
|
|
completed: usize,
|
|
total: usize,
|
|
error: String,
|
|
) -> Self {
|
|
Self {
|
|
workflow_name: name,
|
|
success: false,
|
|
duration,
|
|
steps_completed: completed,
|
|
total_steps: total,
|
|
error_message: Some(error),
|
|
metrics: HashMap::new(),
|
|
order_ids: Vec::new(),
|
|
trades_executed: 0,
|
|
}
|
|
}
|
|
}
|
|
|
|
/// Complete trading workflow orchestrator
|
|
pub struct TradingWorkflow {
|
|
ml_pipeline: Arc<RwLock<MLTestPipeline>>,
|
|
}
|
|
|
|
impl TradingWorkflow {
|
|
pub fn new(
|
|
_database: Arc<TestDatabase>,
|
|
ml_pipeline: Arc<RwLock<MLTestPipeline>>,
|
|
_test_data: Arc<TestDataGenerator>,
|
|
) -> Self {
|
|
Self {
|
|
ml_pipeline,
|
|
}
|
|
}
|
|
|
|
/// Execute complete order lifecycle test
|
|
pub async fn test_order_lifecycle(&self, mut client: TliClient) -> Result<WorkflowTestResult> {
|
|
let start_time = Instant::now();
|
|
let workflow_name = "order_lifecycle_test".to_string();
|
|
|
|
info!("Starting order lifecycle workflow test");
|
|
|
|
let mut steps_completed = 0;
|
|
let total_steps = 8;
|
|
let mut order_ids = Vec::new();
|
|
let mut metrics = HashMap::new();
|
|
|
|
// Step 1: Check system health
|
|
// Note: This would require a monitoring client, for now we'll simulate this step
|
|
info!("System health check - simulated for now (monitoring service needed)");
|
|
steps_completed += 1;
|
|
|
|
// Step 2: Get account information via portfolio summary
|
|
let account_id = "TEST_ACCOUNT_1".to_string();
|
|
if let Some(trading_client) = client.trading() {
|
|
let mut trading_client = trading_client.clone();
|
|
let portfolio_request = GetPortfolioSummaryRequest {
|
|
account_id: account_id.clone(),
|
|
};
|
|
match trading_client
|
|
.get_portfolio_summary(portfolio_request)
|
|
.await
|
|
{
|
|
Ok(portfolio_response) => {
|
|
let portfolio = portfolio_response.into_inner();
|
|
info!("Portfolio value: ${:.2}", portfolio.total_value);
|
|
metrics.insert("initial_balance".to_string(), portfolio.total_value);
|
|
steps_completed += 1;
|
|
},
|
|
Err(e) => {
|
|
return Ok(WorkflowTestResult::failure(
|
|
workflow_name,
|
|
start_time.elapsed(),
|
|
steps_completed,
|
|
total_steps,
|
|
format!("Portfolio info retrieval failed: {}", e),
|
|
));
|
|
},
|
|
}
|
|
}
|
|
|
|
// Step 3: Submit a test order
|
|
let order_request = SubmitOrderRequest {
|
|
symbol: "AAPL".to_string(),
|
|
side: OrderSide::Buy as i32,
|
|
quantity: 100.0,
|
|
order_type: OrderType::Market as i32,
|
|
price: None,
|
|
stop_price: None,
|
|
account_id: account_id.clone(),
|
|
metadata: std::collections::HashMap::new(),
|
|
};
|
|
|
|
let order_id = if let Some(trading_client) = client.trading() {
|
|
let mut trading_client = trading_client.clone();
|
|
match trading_client.submit_order(order_request).await {
|
|
Ok(response) => {
|
|
let submit_response = response.into_inner();
|
|
info!("Order submitted successfully: {}", submit_response.order_id);
|
|
order_ids.push(submit_response.order_id.clone());
|
|
steps_completed += 1;
|
|
submit_response.order_id
|
|
},
|
|
Err(e) => {
|
|
return Ok(WorkflowTestResult::failure(
|
|
workflow_name,
|
|
start_time.elapsed(),
|
|
steps_completed,
|
|
total_steps,
|
|
format!("Order submission error: {}", e),
|
|
));
|
|
},
|
|
}
|
|
} else {
|
|
return Ok(WorkflowTestResult::failure(
|
|
workflow_name,
|
|
start_time.elapsed(),
|
|
steps_completed,
|
|
total_steps,
|
|
"Trading client not available".to_string(),
|
|
));
|
|
};
|
|
|
|
// Step 4: Check order status
|
|
sleep(Duration::from_millis(500)).await; // Allow order to be processed
|
|
|
|
if let Some(trading_client) = client.trading() {
|
|
let mut trading_client = trading_client.clone();
|
|
let status_request = GetOrderStatusRequest {
|
|
order_id: order_id.clone(),
|
|
};
|
|
match trading_client.get_order_status(status_request).await {
|
|
Ok(response) => {
|
|
let order_status = response.into_inner();
|
|
if let Some(order) = order_status.order {
|
|
info!(
|
|
"Order status: {:?}",
|
|
OrderStatus::try_from(order.status).unwrap_or(OrderStatus::Unspecified)
|
|
);
|
|
metrics.insert("order_filled_quantity".to_string(), order.filled_quantity);
|
|
}
|
|
steps_completed += 1;
|
|
},
|
|
Err(e) => {
|
|
return Ok(WorkflowTestResult::failure(
|
|
workflow_name,
|
|
start_time.elapsed(),
|
|
steps_completed,
|
|
total_steps,
|
|
format!("Order status check failed: {}", e),
|
|
));
|
|
},
|
|
}
|
|
}
|
|
|
|
// Step 5: Test risk management (simulated - would require risk service)
|
|
info!("Risk validation - simulated for now (risk service needed)");
|
|
metrics.insert("risk_score".to_string(), 0.75); // Simulated risk score
|
|
steps_completed += 1;
|
|
|
|
// Step 6: Test VaR calculation (simulated - would require risk service)
|
|
info!("VaR calculation - simulated for now (risk service needed)");
|
|
metrics.insert("portfolio_var".to_string(), 25000.0); // Simulated VaR
|
|
steps_completed += 1;
|
|
|
|
// Step 7: Test market data streaming (briefly)
|
|
if let Some(trading_client) = client.trading() {
|
|
let mut trading_client = trading_client.clone();
|
|
let market_data_request = StreamMarketDataRequest {
|
|
symbols: vec!["AAPL".to_string()],
|
|
data_types: vec![MarketDataType::Trade as i32, MarketDataType::Quote as i32],
|
|
};
|
|
match trading_client.stream_market_data(market_data_request).await {
|
|
Ok(response) => {
|
|
let mut stream = response.into_inner();
|
|
info!("Market data stream established");
|
|
|
|
// Collect a few market data events with timeout
|
|
let stream_timeout = Duration::from_secs(5);
|
|
let mut events_received = 0;
|
|
|
|
match timeout(stream_timeout, async {
|
|
while let Some(event) = stream.next().await {
|
|
match event {
|
|
Ok(_market_event) => {
|
|
events_received += 1;
|
|
debug!("Received market data event");
|
|
if events_received >= 3 {
|
|
break;
|
|
}
|
|
},
|
|
Err(e) => {
|
|
warn!("Market data stream error: {}", e);
|
|
break;
|
|
},
|
|
}
|
|
}
|
|
Ok::<(), anyhow::Error>(())
|
|
})
|
|
.await
|
|
{
|
|
Ok(_) => {
|
|
info!(
|
|
"Market data streaming test completed: {} events",
|
|
events_received
|
|
);
|
|
metrics
|
|
.insert("market_data_events".to_string(), events_received as f64);
|
|
steps_completed += 1;
|
|
},
|
|
Err(_) => {
|
|
warn!("Market data streaming timed out, but continuing");
|
|
steps_completed += 1; // Don't fail the entire workflow for streaming timeout
|
|
},
|
|
}
|
|
},
|
|
Err(e) => {
|
|
warn!("Market data subscription failed: {}", e);
|
|
steps_completed += 1; // Don't fail for streaming issues
|
|
},
|
|
}
|
|
}
|
|
|
|
// Step 8: Cancel the test order (cleanup)
|
|
if let Some(trading_client) = client.trading() {
|
|
let mut trading_client = trading_client.clone();
|
|
let cancel_request = CancelOrderRequest {
|
|
order_id: order_id.clone(),
|
|
account_id: account_id.clone(),
|
|
};
|
|
|
|
match trading_client.cancel_order(cancel_request).await {
|
|
Ok(response) => {
|
|
let cancel_response = response.into_inner();
|
|
if cancel_response.success {
|
|
info!("Order cancelled successfully");
|
|
} else {
|
|
info!("Order cancellation not needed (already filled/cancelled)");
|
|
}
|
|
steps_completed += 1;
|
|
},
|
|
Err(e) => {
|
|
warn!("Order cancellation failed: {}", e);
|
|
steps_completed += 1; // Don't fail workflow for cancellation issues
|
|
},
|
|
}
|
|
}
|
|
|
|
let duration = start_time.elapsed();
|
|
info!("Order lifecycle workflow completed in {:?}", duration);
|
|
|
|
let mut result = WorkflowTestResult::success(workflow_name, duration, steps_completed);
|
|
result.metrics = metrics;
|
|
result.order_ids = order_ids;
|
|
result.trades_executed = if steps_completed >= 4 { 1 } else { 0 };
|
|
|
|
Ok(result)
|
|
}
|
|
|
|
/// Test ML-driven trading workflow
|
|
pub async fn test_ml_trading_workflow(
|
|
&self,
|
|
mut client: TliClient,
|
|
) -> Result<WorkflowTestResult> {
|
|
let start_time = Instant::now();
|
|
let workflow_name = "ml_trading_workflow".to_string();
|
|
|
|
info!("Starting ML-driven trading workflow test");
|
|
|
|
let mut steps_completed = 0;
|
|
let total_steps = 6;
|
|
let mut metrics = HashMap::new();
|
|
let mut order_ids = Vec::new();
|
|
|
|
// Step 1: Generate ML predictions
|
|
let test_features: Vec<f64> = (0..50).map(|_| rand::random::<f64>() * 2.0 - 1.0).collect();
|
|
|
|
let ensemble_result = self
|
|
.ml_pipeline
|
|
.write()
|
|
.await
|
|
.test_ensemble_prediction(test_features)
|
|
.await
|
|
.context("ML ensemble prediction failed")?;
|
|
|
|
info!(
|
|
"ML ensemble prediction: {:?} with {:.2}% confidence",
|
|
ensemble_result.prediction,
|
|
ensemble_result.confidence * 100.0
|
|
);
|
|
metrics.insert("ml_confidence".to_string(), ensemble_result.confidence);
|
|
metrics.insert(
|
|
"ml_signal_strength".to_string(),
|
|
ensemble_result.signal_strength,
|
|
);
|
|
steps_completed += 1;
|
|
|
|
// Step 2: Only proceed with trading if ML confidence is high enough
|
|
if ensemble_result.confidence < 0.6 {
|
|
info!(
|
|
"ML confidence too low ({:.2}), skipping trade execution",
|
|
ensemble_result.confidence
|
|
);
|
|
let duration = start_time.elapsed();
|
|
let mut result = WorkflowTestResult::success(workflow_name, duration, steps_completed);
|
|
result.metrics = metrics;
|
|
return Ok(result);
|
|
}
|
|
|
|
// Step 3: Validate ML prediction with risk management
|
|
let (symbol, side, quantity) = match ensemble_result.prediction {
|
|
PredictionType::Buy | PredictionType::StrongBuy => ("AAPL", OrderSide::Buy, 50.0),
|
|
PredictionType::Sell | PredictionType::StrongSell => ("AAPL", OrderSide::Sell, 50.0),
|
|
_ => {
|
|
info!("ML prediction is HOLD, no trade execution needed");
|
|
let duration = start_time.elapsed();
|
|
let mut result =
|
|
WorkflowTestResult::success(workflow_name, duration, steps_completed);
|
|
result.metrics = metrics;
|
|
return Ok(result);
|
|
},
|
|
};
|
|
|
|
// Risk validation simulated (would require risk service)
|
|
info!("Risk management approved ML-driven trade (simulated)");
|
|
metrics.insert("risk_projected_exposure".to_string(), 15000.0);
|
|
steps_completed += 1;
|
|
|
|
// Step 4: Submit ML-driven order
|
|
let order_request = SubmitOrderRequest {
|
|
symbol: symbol.to_string(),
|
|
side: side as i32,
|
|
quantity,
|
|
order_type: OrderType::Limit as i32,
|
|
price: Some(150.0),
|
|
stop_price: None,
|
|
account_id: "TEST_ACCOUNT_1".to_string(),
|
|
metadata: std::collections::HashMap::new(),
|
|
};
|
|
|
|
if let Some(trading_client) = client.trading() {
|
|
let mut trading_client = trading_client.clone();
|
|
match trading_client.submit_order(order_request).await {
|
|
Ok(response) => {
|
|
let submit_response = response.into_inner();
|
|
info!("ML-driven order submitted: {}", submit_response.order_id);
|
|
order_ids.push(submit_response.order_id.clone());
|
|
steps_completed += 1;
|
|
},
|
|
Err(e) => {
|
|
return Ok(WorkflowTestResult::failure(
|
|
workflow_name,
|
|
start_time.elapsed(),
|
|
steps_completed,
|
|
total_steps,
|
|
format!("Order submission error: {}", e),
|
|
));
|
|
},
|
|
}
|
|
}
|
|
|
|
// Step 5: Monitor order execution with streaming
|
|
if let Some(trading_client) = client.trading() {
|
|
let mut trading_client = trading_client.clone();
|
|
let stream_request = StreamOrdersRequest {
|
|
account_id: Some("TEST_ACCOUNT_1".to_string()),
|
|
symbol: None,
|
|
};
|
|
match trading_client.stream_orders(stream_request).await {
|
|
Ok(response) => {
|
|
let mut stream = response.into_inner();
|
|
info!("Order updates stream established");
|
|
|
|
let stream_timeout = Duration::from_secs(10);
|
|
match timeout(stream_timeout, async {
|
|
while let Some(event) = stream.next().await {
|
|
match event {
|
|
Ok(order_event) => {
|
|
info!("Order update: {}", order_event.order_id);
|
|
if let Some(order) = order_event.order {
|
|
if order.filled_quantity > 0.0 {
|
|
metrics.insert(
|
|
"filled_quantity".to_string(),
|
|
order.filled_quantity,
|
|
);
|
|
break;
|
|
}
|
|
}
|
|
},
|
|
Err(e) => {
|
|
warn!("Order update stream error: {}", e);
|
|
break;
|
|
},
|
|
}
|
|
}
|
|
Ok::<(), anyhow::Error>(())
|
|
})
|
|
.await
|
|
{
|
|
Ok(_) => {
|
|
info!("Order monitoring completed");
|
|
steps_completed += 1;
|
|
},
|
|
Err(_) => {
|
|
warn!("Order monitoring timed out");
|
|
steps_completed += 1; // Don't fail for timeout
|
|
},
|
|
}
|
|
},
|
|
Err(e) => {
|
|
warn!("Order updates subscription failed: {}", e);
|
|
steps_completed += 1; // Don't fail workflow
|
|
},
|
|
}
|
|
}
|
|
|
|
// Step 6: Cleanup - cancel any remaining orders
|
|
for order_id in &order_ids {
|
|
if let Some(trading_client) = client.trading() {
|
|
let mut trading_client = trading_client.clone();
|
|
let cancel_request = CancelOrderRequest {
|
|
order_id: order_id.clone(),
|
|
account_id: "test_account".to_string(),
|
|
};
|
|
let _ = trading_client.cancel_order(cancel_request).await; // Best effort cleanup
|
|
}
|
|
}
|
|
steps_completed += 1;
|
|
|
|
let duration = start_time.elapsed();
|
|
info!("ML-driven trading workflow completed in {:?}", duration);
|
|
|
|
let mut result = WorkflowTestResult::success(workflow_name, duration, steps_completed);
|
|
let trades_executed = if metrics.contains_key("filled_quantity") {
|
|
1
|
|
} else {
|
|
0
|
|
};
|
|
result.metrics = metrics;
|
|
result.order_ids = order_ids;
|
|
result.trades_executed = trades_executed;
|
|
|
|
Ok(result)
|
|
}
|
|
|
|
/// Test emergency stop workflow
|
|
pub async fn test_emergency_stop_workflow(
|
|
&self,
|
|
mut client: TliClient,
|
|
) -> Result<WorkflowTestResult> {
|
|
let start_time = Instant::now();
|
|
let workflow_name = "emergency_stop_workflow".to_string();
|
|
|
|
info!("Starting emergency stop workflow test");
|
|
|
|
let mut steps_completed = 0;
|
|
let total_steps = 4;
|
|
let mut metrics = HashMap::new();
|
|
let mut order_ids = Vec::new();
|
|
|
|
// Step 1: Submit some orders to have something to stop
|
|
if let Some(trading_client) = client.trading() {
|
|
let mut trading_client = trading_client.clone();
|
|
for i in 0..3 {
|
|
let order_request = SubmitOrderRequest {
|
|
symbol: "AAPL".to_string(),
|
|
side: if i % 2 == 0 {
|
|
OrderSide::Buy
|
|
} else {
|
|
OrderSide::Sell
|
|
} as i32,
|
|
order_type: OrderType::Limit as i32,
|
|
quantity: 100.0,
|
|
price: Some(if i % 2 == 0 { 145.0 } else { 155.0 }),
|
|
stop_price: None,
|
|
account_id: "test_account".to_string(),
|
|
metadata: {
|
|
let mut map = std::collections::HashMap::new();
|
|
map.insert("time_in_force".to_string(), "DAY".to_string());
|
|
map.insert(
|
|
"client_order_id".to_string(),
|
|
format!("EMERGENCY_TEST_ORDER_{}", i),
|
|
);
|
|
map
|
|
},
|
|
};
|
|
|
|
match trading_client.submit_order(order_request).await {
|
|
Ok(response) => {
|
|
let submit_response = response.into_inner();
|
|
// Check if order was submitted successfully (status indicates success)
|
|
if submit_response.status == OrderStatus::Submitted as i32 {
|
|
order_ids.push(submit_response.order_id);
|
|
}
|
|
},
|
|
Err(e) => {
|
|
warn!("Failed to submit test order {}: {}", i, e);
|
|
},
|
|
}
|
|
|
|
sleep(Duration::from_millis(100)).await;
|
|
}
|
|
|
|
info!(
|
|
"Submitted {} test orders for emergency stop test",
|
|
order_ids.len()
|
|
);
|
|
metrics.insert("orders_submitted".to_string(), order_ids.len() as f64);
|
|
steps_completed += 1;
|
|
}
|
|
|
|
// Step 2: Get portfolio summary as system health check
|
|
if let Some(trading_client) = client.trading() {
|
|
let mut trading_client = trading_client.clone();
|
|
match trading_client
|
|
.get_portfolio_summary(tonic::Request::new(GetPortfolioSummaryRequest {
|
|
account_id: "test_account".to_string(),
|
|
}))
|
|
.await
|
|
{
|
|
Ok(_portfolio) => {
|
|
info!("Pre-emergency portfolio check completed");
|
|
steps_completed += 1;
|
|
},
|
|
Err(e) => {
|
|
return Ok(WorkflowTestResult::failure(
|
|
workflow_name,
|
|
start_time.elapsed(),
|
|
steps_completed,
|
|
total_steps,
|
|
format!("System status check failed: {}", e),
|
|
));
|
|
},
|
|
}
|
|
}
|
|
|
|
// Step 3: Simulate emergency stop by cancelling all orders
|
|
if let Some(trading_client) = client.trading() {
|
|
let mut trading_client = trading_client.clone();
|
|
info!("Simulating emergency stop by cancelling orders");
|
|
let mut orders_cancelled = 0;
|
|
|
|
// Cancel any orders created in previous steps
|
|
for order_id in &order_ids {
|
|
let cancel_request = CancelOrderRequest {
|
|
order_id: order_id.clone(),
|
|
account_id: "test_account".to_string(),
|
|
};
|
|
|
|
if let Ok(_) = trading_client.cancel_order(cancel_request).await {
|
|
orders_cancelled += 1;
|
|
}
|
|
}
|
|
|
|
info!(
|
|
"Emergency simulation completed - Orders cancelled: {}",
|
|
orders_cancelled
|
|
);
|
|
metrics.insert("orders_cancelled".to_string(), orders_cancelled as f64);
|
|
steps_completed += 1;
|
|
}
|
|
|
|
// Step 4: Verify system health after emergency simulation
|
|
sleep(Duration::from_secs(1)).await; // Allow system to process cancellations
|
|
|
|
if let Some(trading_client) = client.trading() {
|
|
let mut trading_client = trading_client.clone();
|
|
match trading_client
|
|
.get_portfolio_summary(tonic::Request::new(GetPortfolioSummaryRequest {
|
|
account_id: "test_account".to_string(),
|
|
}))
|
|
.await
|
|
{
|
|
Ok(_portfolio) => {
|
|
info!("Post-emergency portfolio check completed");
|
|
steps_completed += 1;
|
|
},
|
|
Err(e) => {
|
|
return Ok(WorkflowTestResult::failure(
|
|
workflow_name,
|
|
start_time.elapsed(),
|
|
steps_completed,
|
|
total_steps,
|
|
format!("Post-emergency status check failed: {}", e),
|
|
));
|
|
},
|
|
}
|
|
}
|
|
|
|
let duration = start_time.elapsed();
|
|
info!("Emergency stop workflow completed in {:?}", duration);
|
|
|
|
let mut result = WorkflowTestResult::success(workflow_name, duration, steps_completed);
|
|
result.metrics = metrics;
|
|
result.order_ids = order_ids;
|
|
|
|
Ok(result)
|
|
}
|
|
}
|
|
|
|
/// Backtesting workflow orchestrator
|
|
pub struct BacktestingWorkflow {
|
|
}
|
|
|
|
impl BacktestingWorkflow {
|
|
pub fn new(_database: Arc<TestDatabase>, _test_data: Arc<TestDataGenerator>) -> Self {
|
|
Self {
|
|
}
|
|
}
|
|
|
|
/// Test complete backtesting workflow
|
|
pub async fn test_backtesting_workflow(
|
|
&self,
|
|
mut client: TliClient,
|
|
) -> Result<WorkflowTestResult> {
|
|
let start_time = Instant::now();
|
|
let workflow_name = "backtesting_workflow".to_string();
|
|
|
|
info!("Starting backtesting workflow test");
|
|
|
|
let mut steps_completed = 0;
|
|
let total_steps = 7;
|
|
let mut metrics = HashMap::new();
|
|
|
|
// Step 1: List existing backtests
|
|
let backtest_client = match client.backtesting() {
|
|
Some(client) => client,
|
|
None => {
|
|
return Ok(WorkflowTestResult::failure(
|
|
workflow_name,
|
|
start_time.elapsed(),
|
|
0,
|
|
total_steps,
|
|
"Backtesting client not available".to_string(),
|
|
));
|
|
},
|
|
};
|
|
|
|
match backtest_client
|
|
.list_backtests(ListBacktestsRequest {
|
|
limit: 100,
|
|
offset: 0,
|
|
status_filter: None,
|
|
strategy_name: None,
|
|
})
|
|
.await
|
|
{
|
|
Ok(response) => {
|
|
let list_response = response.into_inner();
|
|
info!("Found {} existing backtests", list_response.backtests.len());
|
|
metrics.insert(
|
|
"existing_backtests".to_string(),
|
|
list_response.backtests.len() as f64,
|
|
);
|
|
steps_completed += 1;
|
|
},
|
|
Err(e) => {
|
|
return Ok(WorkflowTestResult::failure(
|
|
workflow_name,
|
|
start_time.elapsed(),
|
|
steps_completed,
|
|
total_steps,
|
|
format!("Failed to list backtests: {}", e),
|
|
));
|
|
},
|
|
}
|
|
|
|
// Step 2: Start a new backtest
|
|
let backtest_request = StartBacktestRequest {
|
|
strategy_name: "E2E_Test_Strategy".to_string(),
|
|
symbols: vec!["AAPL".to_string(), "GOOGL".to_string()],
|
|
start_date_unix_nanos: (chrono::Utc::now() - chrono::Duration::days(30))
|
|
.timestamp_nanos_opt()
|
|
.unwrap_or(0),
|
|
end_date_unix_nanos: (chrono::Utc::now() - chrono::Duration::days(1))
|
|
.timestamp_nanos_opt()
|
|
.unwrap_or(0),
|
|
initial_capital: 100000.0,
|
|
parameters: HashMap::new(),
|
|
save_results: true,
|
|
description: "E2E test backtest".to_string(),
|
|
};
|
|
|
|
let backtest_id = match backtest_client.start_backtest(backtest_request).await {
|
|
Ok(response) => {
|
|
let backtest_response = response.into_inner();
|
|
if backtest_response.success {
|
|
info!("Backtest started: {}", backtest_response.backtest_id);
|
|
metrics.insert(
|
|
"estimated_duration".to_string(),
|
|
backtest_response.estimated_duration_seconds as f64,
|
|
);
|
|
steps_completed += 1;
|
|
backtest_response.backtest_id
|
|
} else {
|
|
return Ok(WorkflowTestResult::failure(
|
|
workflow_name,
|
|
start_time.elapsed(),
|
|
steps_completed,
|
|
total_steps,
|
|
format!("Backtest start failed: {}", backtest_response.message),
|
|
));
|
|
}
|
|
},
|
|
Err(e) => {
|
|
return Ok(WorkflowTestResult::failure(
|
|
workflow_name,
|
|
start_time.elapsed(),
|
|
steps_completed,
|
|
total_steps,
|
|
format!("Backtest start error: {}", e),
|
|
));
|
|
},
|
|
};
|
|
|
|
// Step 3: Monitor backtest progress
|
|
match backtest_client
|
|
.subscribe_backtest_progress(SubscribeBacktestProgressRequest {
|
|
backtest_id: backtest_id.clone(),
|
|
})
|
|
.await
|
|
{
|
|
Ok(response) => {
|
|
let mut stream = response.into_inner();
|
|
info!("Backtest progress stream established");
|
|
|
|
let monitor_timeout = Duration::from_secs(30);
|
|
let mut progress_updates = 0;
|
|
let mut max_progress: f64 = 0.0;
|
|
|
|
match timeout(monitor_timeout, async {
|
|
while let Some(event) = stream.next().await {
|
|
match event {
|
|
Ok(progress) => {
|
|
progress_updates += 1;
|
|
max_progress = max_progress.max(progress.progress_percentage);
|
|
info!(
|
|
"Backtest progress: {:.1}% ({} trades)",
|
|
progress.progress_percentage, progress.trades_executed
|
|
);
|
|
|
|
if progress.progress_percentage >= 100.0 {
|
|
info!("Backtest completed!");
|
|
break;
|
|
}
|
|
|
|
// For testing purposes, stop after a few updates
|
|
if progress_updates >= 5 {
|
|
break;
|
|
}
|
|
},
|
|
Err(e) => {
|
|
warn!("Backtest progress stream error: {}", e);
|
|
break;
|
|
},
|
|
}
|
|
}
|
|
Ok::<(), anyhow::Error>(())
|
|
})
|
|
.await
|
|
{
|
|
Ok(_) => {
|
|
info!("Backtest progress monitoring completed");
|
|
metrics.insert("progress_updates".to_string(), progress_updates as f64);
|
|
metrics.insert("max_progress".to_string(), max_progress);
|
|
steps_completed += 1;
|
|
},
|
|
Err(_) => {
|
|
warn!("Backtest progress monitoring timed out");
|
|
steps_completed += 1; // Don't fail for timeout
|
|
},
|
|
}
|
|
},
|
|
Err(e) => {
|
|
warn!("Failed to subscribe to backtest progress: {}", e);
|
|
steps_completed += 1; // Don't fail the workflow
|
|
},
|
|
}
|
|
|
|
// Step 4: Check backtest status
|
|
sleep(Duration::from_secs(2)).await; // Allow some processing time
|
|
|
|
match backtest_client
|
|
.get_backtest_status(GetBacktestStatusRequest {
|
|
backtest_id: backtest_id.clone(),
|
|
})
|
|
.await
|
|
{
|
|
Ok(response) => {
|
|
let status = response.into_inner();
|
|
info!(
|
|
"Backtest status: {} ({:.1}% complete)",
|
|
status.status, status.progress_percentage
|
|
);
|
|
metrics.insert("final_progress".to_string(), status.progress_percentage);
|
|
metrics.insert("trades_executed".to_string(), status.trades_executed as f64);
|
|
steps_completed += 1;
|
|
},
|
|
Err(e) => {
|
|
return Ok(WorkflowTestResult::failure(
|
|
workflow_name,
|
|
start_time.elapsed(),
|
|
steps_completed,
|
|
total_steps,
|
|
format!("Backtest status check failed: {}", e),
|
|
));
|
|
},
|
|
}
|
|
|
|
// Step 5: Get backtest results (even if partial)
|
|
match backtest_client
|
|
.get_backtest_results(GetBacktestResultsRequest {
|
|
backtest_id: backtest_id.clone(),
|
|
include_metrics: true,
|
|
include_trades: false,
|
|
})
|
|
.await
|
|
{
|
|
Ok(response) => {
|
|
let results = response.into_inner();
|
|
if let Some(ref metrics_data) = results.metrics {
|
|
info!(
|
|
"Backtest results: {:.2}% return, {:.2} Sharpe ratio",
|
|
metrics_data.total_return * 100.0,
|
|
metrics_data.sharpe_ratio
|
|
);
|
|
metrics.insert("total_return".to_string(), metrics_data.total_return);
|
|
metrics.insert("sharpe_ratio".to_string(), metrics_data.sharpe_ratio);
|
|
metrics.insert("max_drawdown".to_string(), metrics_data.max_drawdown);
|
|
metrics.insert("total_trades".to_string(), metrics_data.total_trades as f64);
|
|
}
|
|
steps_completed += 1;
|
|
},
|
|
Err(e) => {
|
|
warn!("Failed to get backtest results: {}", e);
|
|
steps_completed += 1; // Don't fail if results aren't ready yet
|
|
},
|
|
}
|
|
|
|
// Step 6: Test stopping the backtest (cleanup)
|
|
match backtest_client
|
|
.stop_backtest(StopBacktestRequest {
|
|
backtest_id: backtest_id.clone(),
|
|
save_partial_results: true,
|
|
})
|
|
.await
|
|
{
|
|
Ok(response) => {
|
|
let stop_response = response.into_inner();
|
|
if stop_response.success {
|
|
info!("Backtest stopped successfully");
|
|
} else {
|
|
info!("Backtest stop not needed (already completed)");
|
|
}
|
|
steps_completed += 1;
|
|
},
|
|
Err(e) => {
|
|
warn!("Failed to stop backtest: {}", e);
|
|
steps_completed += 1; // Don't fail for cleanup issues
|
|
},
|
|
}
|
|
|
|
// Step 7: Verify final list of backtests
|
|
match backtest_client
|
|
.list_backtests(ListBacktestsRequest {
|
|
limit: 100,
|
|
offset: 0,
|
|
status_filter: None,
|
|
strategy_name: None,
|
|
})
|
|
.await
|
|
{
|
|
Ok(response) => {
|
|
let list_response = response.into_inner();
|
|
info!("Final backtest count: {}", list_response.backtests.len());
|
|
metrics.insert(
|
|
"final_backtest_count".to_string(),
|
|
list_response.backtests.len() as f64,
|
|
);
|
|
steps_completed += 1;
|
|
},
|
|
Err(e) => {
|
|
warn!("Failed to get final backtest list: {}", e);
|
|
steps_completed += 1; // Don't fail workflow
|
|
},
|
|
}
|
|
|
|
let duration = start_time.elapsed();
|
|
info!("Backtesting workflow completed in {:?}", duration);
|
|
|
|
let mut result = WorkflowTestResult::success(workflow_name, duration, steps_completed);
|
|
result.metrics = metrics;
|
|
|
|
Ok(result)
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
#[test]
|
|
fn test_workflow_test_result_creation() {
|
|
let result = WorkflowTestResult::success("test".to_string(), Duration::from_secs(1), 5);
|
|
assert!(result.success);
|
|
assert_eq!(result.steps_completed, 5);
|
|
assert_eq!(result.total_steps, 5);
|
|
|
|
let failure = WorkflowTestResult::failure(
|
|
"test".to_string(),
|
|
Duration::from_secs(1),
|
|
3,
|
|
5,
|
|
"Test error".to_string(),
|
|
);
|
|
assert!(!failure.success);
|
|
assert_eq!(failure.steps_completed, 3);
|
|
assert_eq!(failure.total_steps, 5);
|
|
assert_eq!(failure.error_message.unwrap(), "Test error");
|
|
}
|
|
}
|