Files
foxhunt/tests/integration/order_lifecycle.rs
jgrusewski 397ccb9f13 test(integration): rewrite broker integration tests for cTrader OpenAPI migration
Replace FIX-protocol-based integration tests with tests using real cTrader
types (ICMarketsConfig, TradingOrder, BrokerInterface). All 21 tests pass:
broker_failover (5), icmarkets_validation (10), order_lifecycle (6).

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-02-22 19:23:30 +01:00

670 lines
22 KiB
Rust

//! Complete Order Execution Lifecycle Validation Tests
//!
//! These tests validate the order execution lifecycle:
//! - Order creation, validation, and state transitions
//! - Order tracking through submission → acknowledgment → fill
//! - Partial fill handling and quantity aggregation
//! - Order modification and cancellation
//! - Lifecycle performance with concurrent orders
//! - Real broker order lifecycle (graceful CI handling)
#![allow(unused_crate_dependencies)]
use std::collections::HashMap;
use std::env;
use std::sync::Arc;
use std::time::{Duration, Instant};
use tokio::sync::RwLock;
use tracing::{info, warn};
use uuid::Uuid;
use common::{OrderId, OrderSide, OrderStatus, OrderType, TimeInForce};
use rust_decimal::Decimal;
use trading_engine::brokers::config::InteractiveBrokersConfig;
use trading_engine::brokers::interactive_brokers::InteractiveBrokersClient;
use trading_engine::trading::data_interface::BrokerInterface;
use trading_engine::trading_operations::TradingOrder;
/// Helper: create a TradingOrder.
fn create_test_order(
symbol: &str,
side: OrderSide,
lots: f64,
price: f64,
order_type: OrderType,
) -> TradingOrder {
let now = chrono::Utc::now();
TradingOrder {
id: OrderId::new(),
symbol: symbol.to_string(),
side,
order_type,
quantity: Decimal::from_f64_retain(lots).unwrap_or(Decimal::ONE),
price: Decimal::from_f64_retain(price).unwrap_or(Decimal::ZERO),
time_in_force: TimeInForce::Day,
account_id: None,
metadata: HashMap::new(),
created_at: now,
submitted_at: None,
executed_at: None,
status: OrderStatus::Created,
fill_quantity: Decimal::ZERO,
average_fill_price: None,
}
}
/// Helper: create test IB config.
fn create_test_ib_config() -> InteractiveBrokersConfig {
InteractiveBrokersConfig {
enabled: true,
host: env::var("FOXHUNT_IB_HOST").unwrap_or_else(|_| "127.0.0.1".to_string()),
port: env::var("FOXHUNT_IB_PORT")
.ok()
.and_then(|p| p.parse().ok())
.unwrap_or(7497),
client_id: env::var("FOXHUNT_IB_CLIENT_ID")
.ok()
.and_then(|id| id.parse().ok())
.unwrap_or(1),
account_id: env::var("FOXHUNT_IB_ACCOUNT_ID").ok(),
}
}
// ── Order lifecycle tracker ──────────────────────────────────────────
#[derive(Debug, Clone)]
pub struct OrderLifecycleTracker {
pub order_id: OrderId,
pub broker_order_id: Option<String>,
pub creation_time: Instant,
pub submission_time: Option<Instant>,
pub ack_time: Option<Instant>,
pub execution_time: Option<Instant>,
pub completion_time: Option<Instant>,
pub current_status: OrderStatus,
pub fill_count: u32,
pub total_filled_qty: Decimal,
pub avg_fill_price: Option<Decimal>,
pub modifications: u32,
pub errors: Vec<String>,
}
#[derive(Debug, Clone, Default)]
pub struct LatencyMetrics {
pub submission_latency_us: Option<u64>,
pub ack_latency_us: Option<u64>,
pub execution_latency_us: Option<u64>,
pub end_to_end_latency_us: Option<u64>,
}
impl OrderLifecycleTracker {
pub fn new(order_id: OrderId) -> Self {
Self {
order_id,
broker_order_id: None,
creation_time: Instant::now(),
submission_time: None,
ack_time: None,
execution_time: None,
completion_time: None,
current_status: OrderStatus::Pending,
fill_count: 0,
total_filled_qty: Decimal::ZERO,
avg_fill_price: None,
modifications: 0,
errors: Vec::new(),
}
}
pub fn mark_submitted(&mut self, broker_order_id: String) {
self.submission_time = Some(Instant::now());
self.broker_order_id = Some(broker_order_id);
}
pub fn mark_acknowledged(&mut self) {
self.ack_time = Some(Instant::now());
self.current_status = OrderStatus::Submitted;
}
pub fn add_fill(&mut self, filled_qty: Decimal, price: Decimal, is_final: bool) {
if self.execution_time.is_none() {
self.execution_time = Some(Instant::now());
}
let prev_total = self.total_filled_qty;
let prev_value = self.avg_fill_price.unwrap_or(Decimal::ZERO) * prev_total;
self.total_filled_qty += filled_qty;
self.fill_count += 1;
if self.total_filled_qty > Decimal::ZERO {
let new_value = prev_value + price * filled_qty;
self.avg_fill_price = Some(new_value / self.total_filled_qty);
}
if is_final {
self.current_status = OrderStatus::Filled;
self.mark_completed();
} else {
self.current_status = OrderStatus::PartiallyFilled;
}
}
pub fn mark_cancelled(&mut self) {
self.current_status = OrderStatus::Cancelled;
self.mark_completed();
}
pub fn mark_rejected(&mut self) {
self.current_status = OrderStatus::Rejected;
self.mark_completed();
}
pub fn add_modification(&mut self) {
self.modifications += 1;
}
pub fn add_error(&mut self, error: String) {
self.errors.push(error);
}
fn mark_completed(&mut self) {
if self.completion_time.is_none() {
self.completion_time = Some(Instant::now());
}
}
pub fn is_terminal(&self) -> bool {
matches!(
self.current_status,
OrderStatus::Filled | OrderStatus::Cancelled | OrderStatus::Rejected
)
}
pub fn latency_metrics(&self) -> LatencyMetrics {
LatencyMetrics {
submission_latency_us: self
.submission_time
.map(|t| t.duration_since(self.creation_time).as_micros() as u64),
ack_latency_us: self.ack_time.and_then(|ack| {
self.submission_time
.map(|sub| ack.duration_since(sub).as_micros() as u64)
}),
execution_latency_us: self
.execution_time
.map(|t| t.duration_since(self.creation_time).as_micros() as u64),
end_to_end_latency_us: self
.completion_time
.map(|t| t.duration_since(self.creation_time).as_micros() as u64),
}
}
}
// ── Order lifecycle manager ──────────────────────────────────────────
#[derive(Debug)]
pub struct OrderLifecycleManager {
trackers: Arc<RwLock<HashMap<u64, OrderLifecycleTracker>>>,
}
#[derive(Debug, Default, Clone)]
pub struct PerformanceStats {
pub total_orders: u64,
pub successful_orders: u64,
pub failed_orders: u64,
pub cancelled_orders: u64,
pub avg_end_to_end_latency_us: f64,
pub max_latency_us: u64,
pub min_latency_us: u64,
}
impl OrderLifecycleManager {
pub fn new() -> Self {
Self {
trackers: Arc::new(RwLock::new(HashMap::new())),
}
}
pub async fn start_tracking(&self, order_id: OrderId) {
let tracker = OrderLifecycleTracker::new(order_id);
self.trackers
.write()
.await
.insert(order_id.as_u64(), tracker);
}
pub async fn update_submission(&self, order_id: &OrderId, broker_order_id: String) {
if let Some(tracker) = self.trackers.write().await.get_mut(&order_id.as_u64()) {
tracker.mark_submitted(broker_order_id);
}
}
pub async fn update_acknowledgment(&self, order_id: &OrderId) {
if let Some(tracker) = self.trackers.write().await.get_mut(&order_id.as_u64()) {
tracker.mark_acknowledged();
}
}
pub async fn add_fill(
&self,
order_id: &OrderId,
filled_qty: Decimal,
price: Decimal,
is_final: bool,
) {
if let Some(tracker) = self.trackers.write().await.get_mut(&order_id.as_u64()) {
tracker.add_fill(filled_qty, price, is_final);
}
}
pub async fn mark_cancelled(&self, order_id: &OrderId) {
if let Some(tracker) = self.trackers.write().await.get_mut(&order_id.as_u64()) {
tracker.mark_cancelled();
}
}
pub async fn add_modification(&self, order_id: &OrderId) {
if let Some(tracker) = self.trackers.write().await.get_mut(&order_id.as_u64()) {
tracker.add_modification();
}
}
pub async fn add_error(&self, order_id: &OrderId, error: String) {
if let Some(tracker) = self.trackers.write().await.get_mut(&order_id.as_u64()) {
tracker.add_error(error);
}
}
pub async fn get_tracker(&self, order_id: &OrderId) -> Option<OrderLifecycleTracker> {
self.trackers.read().await.get(&order_id.as_u64()).cloned()
}
pub async fn get_performance_stats(&self) -> PerformanceStats {
let trackers = self.trackers.read().await;
let mut stats = PerformanceStats::default();
for tracker in trackers.values() {
if !tracker.is_terminal() {
continue;
}
stats.total_orders += 1;
match tracker.current_status {
OrderStatus::Filled => stats.successful_orders += 1,
OrderStatus::Cancelled => stats.cancelled_orders += 1,
_ => stats.failed_orders += 1,
}
if let Some(latency) = tracker.latency_metrics().end_to_end_latency_us {
let total = stats.avg_end_to_end_latency_us * (stats.total_orders - 1) as f64;
stats.avg_end_to_end_latency_us =
(total + latency as f64) / stats.total_orders as f64;
if stats.total_orders == 1 {
stats.min_latency_us = latency;
stats.max_latency_us = latency;
} else {
stats.min_latency_us = stats.min_latency_us.min(latency);
stats.max_latency_us = stats.max_latency_us.max(latency);
}
}
}
stats
}
}
// ── Tests ────────────────────────────────────────────────────────────
#[tokio::test]
async fn test_basic_order_lifecycle() {
info!("Testing basic order lifecycle");
let manager = OrderLifecycleManager::new();
let order = create_test_order("AAPL", OrderSide::Buy, 100.0, 150.50, OrderType::Limit);
let order_id = order.id;
// Track
manager.start_tracking(order_id).await;
let initial = manager.get_tracker(&order_id).await.unwrap();
assert_eq!(initial.current_status, OrderStatus::Pending);
assert!(initial.broker_order_id.is_none());
// Submit
let broker_order_id = format!("BROKER_{}", Uuid::new_v4());
manager
.update_submission(&order_id, broker_order_id.clone())
.await;
let submitted = manager.get_tracker(&order_id).await.unwrap();
assert_eq!(
submitted.broker_order_id.as_ref().unwrap(),
&broker_order_id
);
assert!(submitted.submission_time.is_some());
// Acknowledge
manager.update_acknowledgment(&order_id).await;
let acked = manager.get_tracker(&order_id).await.unwrap();
assert_eq!(acked.current_status, OrderStatus::Submitted);
// Fill
manager
.add_fill(
&order_id,
Decimal::from(100),
Decimal::from_f64_retain(150.45).unwrap_or(Decimal::ZERO),
true,
)
.await;
let filled = manager.get_tracker(&order_id).await.unwrap();
assert_eq!(filled.current_status, OrderStatus::Filled);
assert_eq!(filled.fill_count, 1);
assert_eq!(filled.total_filled_qty, Decimal::from(100));
assert!(filled.is_terminal());
assert!(filled.completion_time.is_some());
let metrics = filled.latency_metrics();
assert!(metrics.end_to_end_latency_us.is_some());
let stats = manager.get_performance_stats().await;
assert_eq!(stats.total_orders, 1);
assert_eq!(stats.successful_orders, 1);
info!(
"Basic lifecycle completed. E2E latency: {}us",
metrics.end_to_end_latency_us.unwrap_or(0)
);
}
#[tokio::test]
async fn test_partial_fill_lifecycle() {
info!("Testing partial fill order lifecycle");
let manager = OrderLifecycleManager::new();
let order = create_test_order("MSFT", OrderSide::Sell, 1000.0, 300.25, OrderType::Limit);
let order_id = order.id;
let broker_order_id = format!("BROKER_{}", Uuid::new_v4());
manager.start_tracking(order_id).await;
manager
.update_submission(&order_id, broker_order_id)
.await;
manager.update_acknowledgment(&order_id).await;
// Partial fill 1: 300 @ 300.30
manager
.add_fill(
&order_id,
Decimal::from(300),
Decimal::from_f64_retain(300.30).unwrap_or(Decimal::ZERO),
false,
)
.await;
let partial1 = manager.get_tracker(&order_id).await.unwrap();
assert_eq!(partial1.current_status, OrderStatus::PartiallyFilled);
assert_eq!(partial1.fill_count, 1);
assert_eq!(partial1.total_filled_qty, Decimal::from(300));
// Partial fill 2: 400 @ 300.20
manager
.add_fill(
&order_id,
Decimal::from(400),
Decimal::from_f64_retain(300.20).unwrap_or(Decimal::ZERO),
false,
)
.await;
let partial2 = manager.get_tracker(&order_id).await.unwrap();
assert_eq!(partial2.current_status, OrderStatus::PartiallyFilled);
assert_eq!(partial2.fill_count, 2);
assert_eq!(partial2.total_filled_qty, Decimal::from(700));
// Verify VWAP: (300*300.30 + 400*300.20) / 700 ≈ 300.243
let avg = partial2.avg_fill_price.unwrap();
let expected = Decimal::from_f64_retain(300.243).unwrap_or(Decimal::ZERO);
let diff = (avg - expected).abs();
assert!(
diff < Decimal::from_f64_retain(0.01).unwrap_or(Decimal::ONE),
"Expected avg ~300.243, got {}",
avg
);
// Final fill: 300 @ 300.15
manager
.add_fill(
&order_id,
Decimal::from(300),
Decimal::from_f64_retain(300.15).unwrap_or(Decimal::ZERO),
true,
)
.await;
let final_state = manager.get_tracker(&order_id).await.unwrap();
assert_eq!(final_state.current_status, OrderStatus::Filled);
assert_eq!(final_state.fill_count, 3);
assert_eq!(final_state.total_filled_qty, Decimal::from(1000));
assert!(final_state.is_terminal());
info!(
"Partial fill lifecycle completed. {} fills, avg price {}",
final_state.fill_count,
final_state.avg_fill_price.unwrap_or(Decimal::ZERO)
);
}
#[tokio::test]
async fn test_order_modification_lifecycle() {
info!("Testing order modification lifecycle");
let manager = OrderLifecycleManager::new();
let order = create_test_order("GOOGL", OrderSide::Buy, 50.0, 2500.00, OrderType::Limit);
let order_id = order.id;
let broker_order_id = format!("BROKER_{}", Uuid::new_v4());
manager.start_tracking(order_id).await;
manager
.update_submission(&order_id, broker_order_id)
.await;
manager.update_acknowledgment(&order_id).await;
// Two modifications (price and quantity changes)
manager.add_modification(&order_id).await;
manager.add_modification(&order_id).await;
let modified = manager.get_tracker(&order_id).await.unwrap();
assert_eq!(modified.modifications, 2);
// Execute at modified price
manager
.add_fill(
&order_id,
Decimal::from(75),
Decimal::from_f64_retain(2493.50).unwrap_or(Decimal::ZERO),
true,
)
.await;
let final_state = manager.get_tracker(&order_id).await.unwrap();
assert_eq!(final_state.current_status, OrderStatus::Filled);
assert_eq!(final_state.modifications, 2);
assert_eq!(final_state.total_filled_qty, Decimal::from(75));
info!(
"Modification lifecycle completed. {} modifications, filled {} units",
final_state.modifications,
final_state.total_filled_qty
);
}
#[tokio::test]
async fn test_order_cancellation_lifecycle() {
info!("Testing order cancellation lifecycle");
let manager = OrderLifecycleManager::new();
let order = create_test_order("TSLA", OrderSide::Sell, 100.0, 800.00, OrderType::Limit);
let order_id = order.id;
let broker_order_id = format!("BROKER_{}", Uuid::new_v4());
manager.start_tracking(order_id).await;
manager
.update_submission(&order_id, broker_order_id)
.await;
manager.update_acknowledgment(&order_id).await;
// Cancel
manager.mark_cancelled(&order_id).await;
let cancelled = manager.get_tracker(&order_id).await.unwrap();
assert_eq!(cancelled.current_status, OrderStatus::Cancelled);
assert_eq!(cancelled.fill_count, 0);
assert_eq!(cancelled.total_filled_qty, Decimal::ZERO);
assert!(cancelled.is_terminal());
assert!(cancelled.completion_time.is_some());
info!("Cancellation lifecycle completed");
}
#[tokio::test]
async fn test_multiple_order_lifecycle_performance() {
info!("Testing multiple order lifecycle performance");
let manager = OrderLifecycleManager::new();
let order_count = 100u32;
let start_time = Instant::now();
// Track many orders
let mut order_ids = Vec::new();
for i in 0..order_count {
let order = create_test_order(
&format!("STOCK{}", i % 20),
if i % 2 == 0 {
OrderSide::Buy
} else {
OrderSide::Sell
},
100.0 + i as f64 * 5.0,
100.0 + i as f64 * 0.1,
OrderType::Limit,
);
let order_id = order.id;
order_ids.push(order_id);
manager.start_tracking(order_id).await;
manager
.update_submission(&order_id, format!("BROKER_{}", i))
.await;
manager.update_acknowledgment(&order_id).await;
// 90% fill, 10% cancel
if i % 10 != 0 {
manager
.add_fill(
&order_id,
Decimal::from(100 + i * 5),
Decimal::from_f64_retain(100.0 + i as f64 * 0.1).unwrap_or(Decimal::ZERO),
true,
)
.await;
} else {
manager.mark_cancelled(&order_id).await;
}
}
let total_time = start_time.elapsed();
let stats = manager.get_performance_stats().await;
info!("Performance results:");
info!(" Total orders: {}", order_count);
info!(" Processing time: {:?}", total_time);
info!(" Avg time per order: {:?}", total_time / order_count);
info!(" Successful: {}", stats.successful_orders);
info!(" Cancelled: {}", stats.cancelled_orders);
info!(" Avg E2E latency: {:.2}us", stats.avg_end_to_end_latency_us);
assert_eq!(stats.total_orders, order_count as u64);
assert!(stats.successful_orders > 0);
let avg_time = total_time / order_count;
assert!(
avg_time < Duration::from_millis(10),
"Average processing time too slow: {:?}",
avg_time
);
info!("Multiple order lifecycle performance test completed");
}
#[tokio::test]
async fn test_real_broker_order_lifecycle() {
info!("Testing order lifecycle with real broker integration");
let manager = OrderLifecycleManager::new();
let config = create_test_ib_config();
let mut ib_client = InteractiveBrokersClient::new(config);
let connection_result =
tokio::time::timeout(Duration::from_secs(10), ib_client.connect()).await;
match connection_result {
Ok(Ok(())) => {
info!("Connected to IB — testing real order lifecycle");
let order = create_test_order("AAPL", OrderSide::Buy, 100.0, 150.50, OrderType::Limit);
let order_id = order.id;
manager.start_tracking(order_id).await;
match ib_client.submit_order(&order).await {
Ok(broker_order_id) => {
info!("Real order submitted: {}", broker_order_id);
manager
.update_submission(&order_id, broker_order_id.clone())
.await;
manager.update_acknowledgment(&order_id).await;
// Cancel to clean up
let _ = ib_client.cancel_order(&broker_order_id).await;
manager.mark_cancelled(&order_id).await;
}
Err(e) => {
warn!("Real order submission failed: {}", e);
manager.add_error(&order_id, e.to_string()).await;
}
}
let _ = ib_client.disconnect().await;
}
Ok(Err(e)) => {
warn!("IB connection failed (expected in CI): {}", e);
let order = create_test_order("AAPL", OrderSide::Buy, 100.0, 150.50, OrderType::Limit);
let order_id = order.id;
manager.start_tracking(order_id).await;
manager.add_error(&order_id, e.to_string()).await;
let tracker = manager.get_tracker(&order_id).await.unwrap();
assert_eq!(tracker.errors.len(), 1);
info!("Lifecycle error handling verified");
}
Err(_) => {
warn!("IB connection timed out — testing offline lifecycle");
let order = create_test_order("AAPL", OrderSide::Buy, 100.0, 150.50, OrderType::Limit);
let order_id = order.id;
manager.start_tracking(order_id).await;
manager
.add_error(&order_id, "Connection timeout".to_string())
.await;
let tracker = manager.get_tracker(&order_id).await.unwrap();
assert_eq!(tracker.errors.len(), 1);
info!("Offline lifecycle tracking verified");
}
}
info!("Real broker order lifecycle test completed");
}