Files
foxhunt/testing/integration/performance_and_stress_tests.rs
jgrusewski 9c3d741a08 refactor: restructure repo — crates/, bin/, testing/ layout
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>
2026-02-25 11:56:00 +01:00

404 lines
14 KiB
Rust

//! Comprehensive Performance and Stress Tests
//!
//! This test suite provides performance benchmarking and stress testing
//! for critical components of the Foxhunt HFT system.
#![allow(unused_crate_dependencies)]
use std::collections::HashMap;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::Instant;
use tokio::sync::RwLock;
// Import common types
use chrono::Utc;
use common::{OrderId, OrderSide, OrderStatus, OrderType, TimeInForce};
use rust_decimal::Decimal;
// Import trading engine modules with correct paths
use trading_engine::lockfree::{HftMessage, LockFreeRingBuffer};
use trading_engine::simd::{AlignedPrices, AlignedVolumes, SimdPriceOps};
use trading_engine::timing::{calibrate_tsc, HardwareTimestamp};
use trading_engine::trading::order_manager::OrderManager;
use trading_engine::trading_operations::TradingOrder;
#[cfg(test)]
mod performance_and_stress_tests {
use super::*;
// ========================================================================
// High-Frequency Trading Performance Tests
// ========================================================================
#[tokio::test]
async fn test_order_processing_latency_target() {
let order_manager = OrderManager::new();
let iterations = 1_000;
// Try to calibrate TSC, but don't fail if it doesn't work
let _ = calibrate_tsc();
let mut latencies_ns = Vec::with_capacity(iterations);
// Warm up
for _ in 0..100 {
let order = create_test_order();
order_manager.add_order(order).await;
}
// Benchmark order processing latency
for i in 0..iterations {
let order = create_test_order_with_id(i);
let start_time = HardwareTimestamp::now();
order_manager.add_order(order).await;
let end_time = HardwareTimestamp::now();
let latency_ns = end_time.latency_ns(&start_time);
latencies_ns.push(latency_ns);
}
// Calculate percentiles
latencies_ns.sort_unstable();
let p50 = latencies_ns[latencies_ns.len() / 2];
let p95 = latencies_ns[(latencies_ns.len() as f64 * 0.95) as usize];
let p99 = latencies_ns[(latencies_ns.len() as f64 * 0.99) as usize];
println!("Order Processing Latency Results:");
println!(" P50: {}ns ({:.2}μs)", p50, p50 as f64 / 1000.0);
println!(" P95: {}ns ({:.2}μs)", p95, p95 as f64 / 1000.0);
println!(" P99: {}ns ({:.2}μs)", p99, p99 as f64 / 1000.0);
// Verify sub-millisecond performance (relaxed for test environment)
assert!(p95 < 10_000_000, "P95 latency {}ns exceeds 10ms", p95);
}
#[tokio::test]
async fn test_simd_price_calculations_performance() {
if !std::arch::is_x86_feature_detected!("avx2") {
println!("Skipping SIMD test - AVX2 not available");
return;
}
const ARRAY_SIZE: usize = 10_000;
const ITERATIONS: usize = 100;
// Generate test price data
let prices: Vec<f64> = (0..ARRAY_SIZE).map(|i| 100.0 + (i as f64 * 0.01)).collect();
let volumes: Vec<f64> = (0..ARRAY_SIZE)
.map(|i| 1000.0 + (i as f64 * 10.0))
.collect();
let aligned_prices = AlignedPrices::from_slice(&prices);
let aligned_volumes = AlignedVolumes::from_slice(&volumes);
// SIMD performance test
let simd_start = Instant::now();
let mut simd_results = Vec::new();
// SAFETY: AVX2 support verified above
unsafe {
let simd_ops = SimdPriceOps::new();
for _ in 0..ITERATIONS {
let vwap = simd_ops.calculate_vwap_aligned(&aligned_prices, &aligned_volumes);
simd_results.push(vwap);
}
}
let simd_time = simd_start.elapsed();
// Scalar performance test (for comparison)
let scalar_start = Instant::now();
let mut scalar_results = Vec::new();
for _ in 0..ITERATIONS {
let total_pv: f64 = prices.iter().zip(volumes.iter()).map(|(p, v)| p * v).sum();
let total_volume: f64 = volumes.iter().sum();
let vwap = if total_volume > 0.0 {
total_pv / total_volume
} else {
0.0
};
scalar_results.push(vwap);
}
let scalar_time = scalar_start.elapsed();
// Verify results are equivalent
assert_eq!(simd_results.len(), scalar_results.len());
for (simd, scalar) in simd_results.into_iter().zip(scalar_results.into_iter()) {
assert!(
(simd - scalar).abs() < 1e-10,
"SIMD and scalar results differ: {} vs {}",
simd,
scalar
);
}
let simd_throughput = (ITERATIONS * ARRAY_SIZE) as f64 / simd_time.as_secs_f64();
let scalar_throughput = (ITERATIONS * ARRAY_SIZE) as f64 / scalar_time.as_secs_f64();
let speedup = scalar_time.as_secs_f64() / simd_time.as_secs_f64();
println!("SIMD Performance Comparison:");
println!(" SIMD time: {:?}", simd_time);
println!(" Scalar time: {:?}", scalar_time);
println!(" SIMD throughput: {:.0} ops/sec", simd_throughput);
println!(" Scalar throughput: {:.0} ops/sec", scalar_throughput);
println!(" Speedup ratio: {:.2}x", speedup);
// SIMD should provide some speedup
assert!(simd_time <= scalar_time, "SIMD not faster than scalar");
}
#[tokio::test]
async fn test_lock_free_structures_performance() {
const NUM_MESSAGES: usize = 100_000;
let ring_buffer = Arc::new(
LockFreeRingBuffer::<HftMessage>::new(1_000_000).expect("Failed to create ring buffer"),
);
let start_time = Instant::now();
let processed_count = Arc::new(AtomicU64::new(0));
// Spawn producer task
let buffer_clone = Arc::clone(&ring_buffer);
let producer_handle = tokio::spawn(async move {
let mut local_count = 0;
for msg_id in 0..NUM_MESSAGES {
let message = HftMessage::new(1, [msg_id as u64, 0, 0, 0, 0, 0, 0, 0]);
while buffer_clone.try_push(message).is_err() {
tokio::task::yield_now().await; // Back pressure
}
local_count += 1;
}
local_count
});
// Spawn consumer task
let buffer_clone = Arc::clone(&ring_buffer);
let counter_clone = Arc::clone(&processed_count);
let consumer_handle = tokio::spawn(async move {
let mut local_count = 0;
loop {
if let Some(_event) = buffer_clone.try_pop() {
local_count += 1;
counter_clone.fetch_add(1, Ordering::Relaxed);
} else {
tokio::task::yield_now().await;
}
// Check if we've processed all messages
if counter_clone.load(Ordering::Relaxed) >= NUM_MESSAGES as u64 {
break;
}
}
local_count
});
// Wait for completion
let total_produced = producer_handle.await.expect("Producer task failed");
let total_consumed = consumer_handle.await.expect("Consumer task failed");
let total_time = start_time.elapsed();
let throughput = NUM_MESSAGES as f64 / total_time.as_secs_f64();
println!("Lock-Free Ring Buffer Performance:");
println!(" Total messages: {}", NUM_MESSAGES);
println!(
" Produced: {}, Consumed: {}",
total_produced, total_consumed
);
println!(" Processing time: {:?}", total_time);
println!(" Throughput: {:.0} messages/sec", throughput);
assert_eq!(total_produced, NUM_MESSAGES);
assert_eq!(total_consumed, NUM_MESSAGES);
assert!(
throughput > 100_000.0,
"Lock-free throughput {:.0} below 100K messages/s",
throughput
);
}
// ========================================================================
// Stress Testing - System Under Load
// ========================================================================
#[tokio::test]
async fn test_concurrent_order_processing_stress() {
let order_manager = Arc::new(RwLock::new(OrderManager::new()));
let concurrent_traders = 10;
let orders_per_trader = 100;
let total_orders = concurrent_traders * orders_per_trader;
let start_time = Instant::now();
let success_counter = Arc::new(AtomicU64::new(0));
// Launch concurrent trading sessions
let mut trader_handles = Vec::new();
for trader_id in 0..concurrent_traders {
let manager_clone = Arc::clone(&order_manager);
let success_clone = Arc::clone(&success_counter);
let handle = tokio::spawn(async move {
let mut local_success = 0;
for order_id in 0..orders_per_trader {
let order = TradingOrder {
id: OrderId::from(format!("TRADER{:03}_{:06}", trader_id, order_id)),
symbol: format!("SYMBOL{:02}", order_id % 10),
side: if order_id % 2 == 0 {
OrderSide::Buy
} else {
OrderSide::Sell
},
order_type: OrderType::Limit,
quantity: Decimal::from((order_id + 1) * 1000),
price: Decimal::new(10000 + order_id as i64, 4), // e.g. 1.0001
time_in_force: TimeInForce::GoodTillCancel,
account_id: None,
metadata: HashMap::new(),
created_at: Utc::now(),
submitted_at: None,
executed_at: None,
status: OrderStatus::New,
fill_quantity: Decimal::ZERO,
average_fill_price: Some(Decimal::ZERO),
};
let mgr = manager_clone.write().await;
mgr.add_order(order).await;
local_success += 1;
success_clone.fetch_add(1, Ordering::Relaxed);
// Simulate realistic trading pace
if order_id % 10 == 0 {
tokio::task::yield_now().await;
}
}
local_success
});
trader_handles.push(handle);
}
// Wait for all traders to complete
let mut total_success = 0;
for handle in trader_handles {
let success = handle.await.expect("Trader task failed");
total_success += success;
}
let total_time = start_time.elapsed();
let throughput = total_success as f64 / total_time.as_secs_f64();
println!("Concurrent Order Processing Stress Test:");
println!(" Concurrent traders: {}", concurrent_traders);
println!(" Orders per trader: {}", orders_per_trader);
println!(" Total orders: {}", total_orders);
println!(" Successful orders: {}", total_success);
println!(" Total time: {:?}", total_time);
println!(" Throughput: {:.0} orders/sec", throughput);
// All orders should succeed
assert_eq!(total_success, total_orders);
}
#[tokio::test]
async fn test_memory_pressure_handling() {
// Test system behavior under memory pressure
let large_allocation_size = 100_000; // 100K elements
let num_allocations = 10;
let mut allocations = Vec::new();
println!("Starting memory pressure test...");
// Gradually increase memory pressure
for i in 0..num_allocations {
let allocation: Vec<f64> = (0..large_allocation_size)
.map(|j| (i * large_allocation_size + j) as f64)
.collect();
allocations.push(allocation);
// Test system responsiveness under memory pressure
if i % 5 == 0 {
// Verify system can still process orders
let order_manager = OrderManager::new();
let test_order = create_test_order();
let start_time = Instant::now();
order_manager.add_order(test_order).await;
let latency = start_time.elapsed();
assert!(
latency.as_millis() < 1000,
"Order latency {}ms too high under memory pressure",
latency.as_millis()
);
}
}
// Clean up
allocations.clear();
// Force some async yields
for _ in 0..10 {
tokio::task::yield_now().await;
}
println!("Memory pressure test completed successfully");
}
}
// ============================================================================
// Test Utilities and Helper Functions
// ============================================================================
fn create_test_order() -> TradingOrder {
TradingOrder {
id: OrderId::from(uuid::Uuid::new_v4().to_string()),
symbol: "EURUSD".to_string(),
side: OrderSide::Buy,
order_type: OrderType::Limit,
quantity: Decimal::from(10000),
price: Decimal::new(12345, 4), // 1.2345
time_in_force: TimeInForce::GoodTillCancel,
account_id: None,
metadata: HashMap::new(),
created_at: Utc::now(),
submitted_at: None,
executed_at: None,
status: OrderStatus::New,
fill_quantity: Decimal::ZERO,
average_fill_price: Some(Decimal::ZERO),
}
}
fn create_test_order_with_id(id: usize) -> TradingOrder {
TradingOrder {
id: OrderId::from(format!("TEST_ORDER_{:06}", id)),
symbol: "EURUSD".to_string(),
side: if id % 2 == 0 {
OrderSide::Buy
} else {
OrderSide::Sell
},
order_type: OrderType::Limit,
quantity: Decimal::from((id + 1) * 1000),
price: Decimal::new(12345 + id as i64, 4), // 1.2345 + small increment
time_in_force: TimeInForce::GoodTillCancel,
account_id: None,
metadata: HashMap::new(),
created_at: Utc::now(),
submitted_at: None,
executed_at: None,
status: OrderStatus::New,
fill_quantity: Decimal::ZERO,
average_fill_price: Some(Decimal::ZERO),
}
}