Files
foxhunt/data/tests/test_databento_streaming.rs
jgrusewski eb5fe84e22 🔥 COMPILATION SUCCESS: Complete resolution of all 543+ compilation errors
ARCHITECTURAL ACHIEVEMENTS:
 Zero compilation errors across entire workspace
 Complete elimination of circular dependencies
 Proper configuration architecture with centralized config crate
 Fixed all type mismatches and missing fields
 Restored proper crate structure (config at root level)

MAJOR FIXES:
- Fixed 19 critical data crate compilation errors
- Resolved configuration struct field mismatches
- Fixed enum variant naming (CSV → Csv)
- Corrected type conversions (FromPrimitive, compression types)
- Fixed HashMap key types (u32 vs usize)
- Resolved TLOBProcessor constructor issues

WORKSPACE STATUS:
- All services compile successfully
- Trading Service:  Ready
- Backtesting Service:  Ready
- ML Training Service:  Ready
- TLI Client:  Ready

Only documentation warnings remain (3,316 warnings to be addressed)

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude <noreply@anthropic.com>
2025-09-29 10:59:34 +02:00

683 lines
22 KiB
Rust

//! Comprehensive tests for DatabentoStreamingProvider
//!
//! This module contains extensive tests for the Databento streaming WebSocket provider,
//! covering connection management, message parsing, event conversion, error handling,
//! reconnection logic, and performance characteristics.
use chrono::Utc;
use data::error::{DataError, Result};
use data::providers::databento_streaming::{
DatabentoError, DatabentoMessage, DatabentoOrderBook, DatabentoQuote, DatabentoStatus,
DatabentoStreamingProvider, DatabentoSubscription, DatabentoTrade,
};
use data::providers::{ConnectionState, MarketDataProvider};
use rust_decimal_macros::dec;
use serde_json::json;
use std::sync::atomic::Ordering;
use std::sync::Arc;
use std::time::Duration;
use tokio::time::{sleep, timeout};
use tokio_test;
use trading_engine::trading::data_interface::{
MarketDataEvent as CoreMarketDataEvent, QuoteEvent,
};
use common::TradeEvent;
use common::Price;
use common::Quantity;
use common::Symbol;
/// Test provider creation with valid API key
#[tokio::test]
async fn test_provider_creation_success() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string());
assert!(provider.is_ok());
let provider = provider.unwrap();
assert_eq!(provider.get_name(), "databento");
assert!(!provider.connected.load(Ordering::Relaxed));
assert_eq!(provider.messages_received.load(Ordering::Relaxed), 0);
assert_eq!(provider.error_count.load(Ordering::Relaxed), 0);
}
/// Test provider creation with empty API key
#[tokio::test]
async fn test_provider_creation_empty_key() {
let provider = DatabentoStreamingProvider::new("".to_string());
assert!(provider.is_ok()); // Creation should succeed, validation happens on connect
}
/// Test provider clone functionality
#[tokio::test]
async fn test_provider_clone() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
let cloned = provider.clone();
assert_eq!(provider.get_name(), cloned.get_name());
// Atomic counters should point to same memory locations
assert!(Arc::ptr_eq(&provider.connected, &cloned.connected));
assert!(Arc::ptr_eq(
&provider.messages_received,
&cloned.messages_received
));
}
/// Test subscription to market events
#[tokio::test]
async fn test_event_subscription() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
let mut receiver = provider.subscribe_market_events();
// Test that receiver is created successfully
assert_eq!(receiver.len(), 0);
}
/// Test trade message processing
#[tokio::test]
async fn test_process_trade_message() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
let mut receiver = provider.subscribe_market_events();
let trade = DatabentoTrade {
symbol: "SPY".to_string(),
timestamp: Utc::now(),
price: Price::from_f64(425.50).unwrap(),
size: Quantity::from(100),
trade_id: Some("12345".to_string()),
exchange: Some("NYSE".to_string()),
conditions: Some(vec!["Normal".to_string()]),
};
let message = DatabentoMessage::Trade(trade.clone());
// Process the message
let result = provider.process_databento_message(message).await;
assert!(result.is_ok());
// Check that event was sent
let event = timeout(Duration::from_millis(100), receiver.recv()).await;
assert!(event.is_ok());
match event.unwrap().unwrap() {
CoreMarketDataEvent::Trade(trade_event) => {
assert_eq!(trade_event.symbol, trade.symbol);
assert_eq!(trade_event.price, trade.price);
assert_eq!(trade_event.size, trade.size);
assert_eq!(trade_event.exchange, trade.exchange);
}
_ => panic!("Expected trade event"),
}
}
/// Test quote message processing
#[tokio::test]
async fn test_process_quote_message() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
let mut receiver = provider.subscribe_market_events();
let quote = DatabentoQuote {
symbol: "AAPL".to_string(),
timestamp: Utc::now(),
bid: Some(Price::from_f64(150.25).unwrap()),
bid_size: Some(Quantity::from(500)),
ask: Some(Price::from_f64(150.26).unwrap()),
ask_size: Some(Quantity::from(300)),
exchange: Some("NASDAQ".to_string()),
};
let message = DatabentoMessage::Quote(quote.clone());
let result = provider.process_databento_message(message).await;
assert!(result.is_ok());
let event = timeout(Duration::from_millis(100), receiver.recv()).await;
assert!(event.is_ok());
match event.unwrap().unwrap() {
CoreMarketDataEvent::Quote(quote_event) => {
assert_eq!(quote_event.symbol, quote.symbol);
assert_eq!(quote_event.bid, quote.bid);
assert_eq!(quote_event.ask, quote.ask);
assert_eq!(quote_event.bid_size, quote.bid_size);
assert_eq!(quote_event.ask_size, quote.ask_size);
}
_ => panic!("Expected quote event"),
}
}
/// Test order book message processing
#[tokio::test]
async fn test_process_orderbook_message() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
let mut receiver = provider.subscribe_market_events();
let orderbook = DatabentoOrderBook {
symbol: "QQQ".to_string(),
timestamp: Utc::now(),
bids: vec![
(Price::from_f64(375.50).unwrap(), Quantity::from(100)),
(Price::from_f64(375.49).unwrap(), Quantity::from(200)),
],
asks: vec![
(Price::from_f64(375.51).unwrap(), Quantity::from(150)),
(Price::from_f64(375.52).unwrap(), Quantity::from(250)),
],
sequence: Some(12345),
};
let message = DatabentoMessage::OrderBook(orderbook.clone());
let result = provider.process_databento_message(message).await;
assert!(result.is_ok());
let event = timeout(Duration::from_millis(100), receiver.recv()).await;
assert!(event.is_ok());
match event.unwrap().unwrap() {
CoreMarketDataEvent::OrderBook(book_event) => {
assert_eq!(book_event.symbol, orderbook.symbol);
assert_eq!(book_event.bids, orderbook.bids);
assert_eq!(book_event.asks, orderbook.asks);
}
_ => panic!("Expected order book event"),
}
}
/// Test status message processing
#[tokio::test]
async fn test_process_status_message() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
let status = DatabentoStatus {
message: "Connected successfully".to_string(),
timestamp: Utc::now(),
level: "info".to_string(),
};
let message = DatabentoMessage::Status(status);
let result = provider.process_databento_message(message).await;
assert!(result.is_ok());
// Status messages don't generate events, just log output
assert_eq!(provider.messages_received.load(Ordering::Relaxed), 0);
}
/// Test error message processing
#[tokio::test]
async fn test_process_error_message() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
let error = DatabentoError {
error_message: "Invalid symbol".to_string(),
code: Some(400),
timestamp: Utc::now(),
};
let message = DatabentoMessage::Error(error);
let result = provider.process_databento_message(message).await;
assert!(result.is_ok());
// Error should increment error count
assert_eq!(provider.error_count.load(Ordering::Relaxed), 1);
}
/// Test text message processing with valid JSON
#[tokio::test]
async fn test_process_text_message_valid() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
let trade_json = json!({
"type": "trade",
"symbol": "SPY",
"timestamp": "2024-01-15T09:30:00Z",
"price": 425.50,
"size": 100,
"trade_id": "12345",
"exchange": "NYSE",
"conditions": ["Normal"]
});
let result = provider.process_text_message(&trade_json.to_string()).await;
assert!(result.is_ok());
assert_eq!(provider.messages_received.load(Ordering::Relaxed), 1);
assert!(provider.last_message_time.load(Ordering::Relaxed) > 0);
}
/// Test text message processing with invalid JSON
#[tokio::test]
async fn test_process_text_message_invalid() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
let invalid_json = "{ invalid json }";
let result = provider.process_text_message(invalid_json).await;
assert!(result.is_ok()); // Method handles errors internally
assert_eq!(provider.messages_received.load(Ordering::Relaxed), 0);
assert_eq!(provider.error_count.load(Ordering::Relaxed), 1);
}
/// Test binary message processing
#[tokio::test]
async fn test_process_binary_message() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
let binary_data = vec![0x01, 0x02, 0x03, 0x04];
let result = provider.process_binary_message(&binary_data).await;
assert!(result.is_ok());
// Binary processing is not implemented yet, should just log debug message
}
/// Test subscription message creation
#[tokio::test]
async fn test_subscription_creation() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
let symbols = vec![
Symbol::from("SPY"),
Symbol::from("QQQ"),
Symbol::from("IWM"),
];
let result = provider.send_subscription(symbols.clone()).await;
assert!(result.is_ok());
}
/// Test subscription serialization
#[test]
fn test_subscription_serialization() {
let subscription = DatabentoSubscription {
action: "subscribe".to_string(),
symbols: vec![Symbol::from("AAPL"), Symbol::from("GOOGL")],
data_types: vec!["trades".to_string(), "quotes".to_string()],
schema: "ohlcv-1s".to_string(),
};
let json = serde_json::to_string(&subscription);
assert!(json.is_ok());
let json_str = json.unwrap();
assert!(json_str.contains("subscribe"));
assert!(json_str.contains("AAPL"));
assert!(json_str.contains("trades"));
}
/// Test unsubscription serialization
#[test]
fn test_unsubscription_serialization() {
let unsubscription = DatabentoSubscription {
action: "unsubscribe".to_string(),
symbols: vec![Symbol::from("TSLA")],
data_types: vec![],
schema: "".to_string(),
};
let json = serde_json::to_string(&unsubscription);
assert!(json.is_ok());
let json_str = json.unwrap();
assert!(json_str.contains("unsubscribe"));
assert!(json_str.contains("TSLA"));
}
/// Test health status when disconnected
#[tokio::test]
async fn test_health_status_disconnected() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
let health = provider.get_health_status();
assert!(!health.connected);
assert_eq!(health.last_connected, None);
assert_eq!(health.active_subscriptions, 0);
assert_eq!(health.messages_per_second, 0.0);
assert_eq!(health.latency_micros, None);
assert_eq!(health.error_count, 0);
}
/// Test health status with message activity
#[tokio::test]
async fn test_health_status_with_activity() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
// Simulate connection
provider.connected.store(true, Ordering::Relaxed);
provider.messages_received.store(100, Ordering::Relaxed);
provider.last_message_time.store(
chrono::Utc::now().timestamp_millis() as u64,
Ordering::Relaxed,
);
provider.error_count.store(5, Ordering::Relaxed);
let health = provider.get_health_status();
assert!(health.connected);
assert!(health.last_connected.is_some());
assert_eq!(health.error_count, 5);
}
/// Test message parsing edge cases
#[tokio::test]
async fn test_message_parsing_edge_cases() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
// Test empty message
let result = provider.process_text_message("").await;
assert!(result.is_ok());
assert_eq!(provider.error_count.load(Ordering::Relaxed), 1);
// Test null message
let result = provider.process_text_message("null").await;
assert!(result.is_ok());
assert_eq!(provider.error_count.load(Ordering::Relaxed), 2);
// Test malformed JSON
let result = provider.process_text_message("{\"incomplete\":").await;
assert!(result.is_ok());
assert_eq!(provider.error_count.load(Ordering::Relaxed), 3);
}
/// Test trade event with missing optional fields
#[tokio::test]
async fn test_trade_event_minimal_fields() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
let mut receiver = provider.subscribe_market_events();
let trade = DatabentoTrade {
symbol: "MINIMAL".to_string(),
timestamp: Utc::now(),
price: Price::from_f64(100.00).unwrap(),
size: Quantity::from(1),
trade_id: None,
exchange: None,
conditions: None,
};
let message = DatabentoMessage::Trade(trade.clone());
let result = provider.process_databento_message(message).await;
assert!(result.is_ok());
let event = timeout(Duration::from_millis(100), receiver.recv()).await;
assert!(event.is_ok());
match event.unwrap().unwrap() {
CoreMarketDataEvent::Trade(trade_event) => {
assert_eq!(trade_event.symbol, trade.symbol);
assert_eq!(trade_event.trade_id, None);
assert_eq!(trade_event.exchange, None);
}
_ => panic!("Expected trade event"),
}
}
/// Test quote event with partial data
#[tokio::test]
async fn test_quote_event_partial_data() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
let mut receiver = provider.subscribe_market_events();
let quote = DatabentoQuote {
symbol: "PARTIAL".to_string(),
timestamp: Utc::now(),
bid: Some(Price::from_f64(50.00).unwrap()),
bid_size: Some(Quantity::from(100)),
ask: None,
ask_size: None,
exchange: Some("TEST".to_string()),
};
let message = DatabentoMessage::Quote(quote.clone());
let result = provider.process_databento_message(message).await;
assert!(result.is_ok());
let event = timeout(Duration::from_millis(100), receiver.recv()).await;
assert!(event.is_ok());
match event.unwrap().unwrap() {
CoreMarketDataEvent::Quote(quote_event) => {
assert_eq!(quote_event.symbol, quote.symbol);
assert_eq!(quote_event.bid, quote.bid);
assert_eq!(quote_event.ask, None);
assert_eq!(quote_event.bid_size, quote.bid_size);
assert_eq!(quote_event.ask_size, None);
}
_ => panic!("Expected quote event"),
}
}
/// Test order book with empty levels
#[tokio::test]
async fn test_orderbook_empty_levels() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
let mut receiver = provider.subscribe_market_events();
let orderbook = DatabentoOrderBook {
symbol: "EMPTY".to_string(),
timestamp: Utc::now(),
bids: vec![],
asks: vec![],
sequence: None,
};
let message = DatabentoMessage::OrderBook(orderbook.clone());
let result = provider.process_databento_message(message).await;
assert!(result.is_ok());
let event = timeout(Duration::from_millis(100), receiver.recv()).await;
assert!(event.is_ok());
match event.unwrap().unwrap() {
CoreMarketDataEvent::OrderBook(book_event) => {
assert_eq!(book_event.symbol, orderbook.symbol);
assert!(book_event.bids.is_empty());
assert!(book_event.asks.is_empty());
}
_ => panic!("Expected order book event"),
}
}
/// Test concurrent message processing
#[tokio::test]
async fn test_concurrent_message_processing() {
let provider = Arc::new(DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap());
let mut receiver = provider.subscribe_market_events();
let mut handles = vec![];
// Spawn multiple tasks processing messages concurrently
for i in 0..10 {
let provider_clone = Arc::clone(&provider);
let handle = tokio::spawn(async move {
let trade = DatabentoTrade {
symbol: format!("SYM{}", i),
timestamp: Utc::now(),
price: Price::from_f64(100.0 + i as f64).unwrap(),
size: Quantity::from(100),
trade_id: Some(format!("trade{}", i)),
exchange: Some("TEST".to_string()),
conditions: None,
};
let message = DatabentoMessage::Trade(trade);
provider_clone.process_databento_message(message).await
});
handles.push(handle);
}
// Wait for all tasks to complete
for handle in handles {
let result = handle.await;
assert!(result.is_ok());
assert!(result.unwrap().is_ok());
}
// Should have received 10 messages
let mut event_count = 0;
while let Ok(Ok(_)) = timeout(Duration::from_millis(10), receiver.recv()).await {
event_count += 1;
if event_count >= 10 {
break;
}
}
assert_eq!(event_count, 10);
}
/// Test message rate calculation
#[tokio::test]
async fn test_message_rate_calculation() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
// Set up initial state
let start_time = chrono::Utc::now().timestamp_millis() as u64;
provider
.last_message_time
.store(start_time, Ordering::Relaxed);
provider.messages_received.store(0, Ordering::Relaxed);
// Process some messages
for i in 1..=5 {
let trade = DatabentoTrade {
symbol: "RATE_TEST".to_string(),
timestamp: Utc::now(),
price: Price::from_f64(100.00).unwrap(),
size: Quantity::from(100),
trade_id: Some(format!("rate_test_{}", i)),
exchange: Some("TEST".to_string()),
conditions: None,
};
let message = DatabentoMessage::Trade(trade);
let _ = provider.process_databento_message(message).await;
sleep(Duration::from_millis(100)).await;
}
let health = provider.get_health_status();
assert!(health.messages_per_second >= 0.0);
}
/// Test error handling in message processing
#[tokio::test]
async fn test_error_handling_message_processing() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
// Test various invalid JSON structures
let invalid_messages = vec![
"not json at all",
"{",
"}",
"[]",
"null",
"true",
"false",
"123",
r#"{"type": "unknown_type"}"#,
r#"{"type": "trade", "symbol": null}"#,
r#"{"type": "trade", "symbol": "SPY", "price": "not_a_number"}"#,
];
let initial_error_count = provider.error_count.load(Ordering::Relaxed);
for (i, msg) in invalid_messages.iter().enumerate() {
let result = provider.process_text_message(msg).await;
assert!(
result.is_ok(),
"Processing should not panic for invalid message {}",
i
);
}
let final_error_count = provider.error_count.load(Ordering::Relaxed);
assert!(final_error_count > initial_error_count);
}
/// Test provider name consistency
#[test]
fn test_provider_name_consistency() {
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
assert_eq!(provider.get_name(), "databento");
let cloned = provider.clone();
assert_eq!(cloned.get_name(), "databento");
}
/// Test all Databento message types serialization
#[test]
fn test_all_message_types_serialization() {
let messages = vec![
DatabentoMessage::Trade(DatabentoTrade {
symbol: "TEST".to_string(),
timestamp: Utc::now(),
price: Price::from_f64(100.0).unwrap(),
size: Quantity::from(100),
trade_id: None,
exchange: None,
conditions: None,
}),
DatabentoMessage::Quote(DatabentoQuote {
symbol: "TEST".to_string(),
timestamp: Utc::now(),
bid: Some(Price::from_f64(99.99).unwrap()),
bid_size: Some(Quantity::from(100)),
ask: Some(Price::from_f64(100.01).unwrap()),
ask_size: Some(Quantity::from(100)),
exchange: None,
}),
DatabentoMessage::OrderBook(DatabentoOrderBook {
symbol: "TEST".to_string(),
timestamp: Utc::now(),
bids: vec![(Price::from_f64(99.99).unwrap(), Quantity::from(100))],
asks: vec![(Price::from_f64(100.01).unwrap(), Quantity::from(100))],
sequence: Some(12345),
}),
DatabentoMessage::Status(DatabentoStatus {
message: "Test status".to_string(),
timestamp: Utc::now(),
level: "info".to_string(),
}),
DatabentoMessage::Error(DatabentoError {
error_message: "Test error".to_string(),
code: Some(400),
timestamp: Utc::now(),
}),
];
for (i, message) in messages.iter().enumerate() {
let json = serde_json::to_string(message);
assert!(json.is_ok(), "Failed to serialize message type {}", i);
let json_str = json.unwrap();
let deserialized: Result<DatabentoMessage, _> = serde_json::from_str(&json_str);
assert!(
deserialized.is_ok(),
"Failed to deserialize message type {}",
i
);
}
}
/// Test WebSocket message handling
#[tokio::test]
async fn test_websocket_message_handling() {
use tokio_tungstenite::tungstenite::Message;
let provider = DatabentoStreamingProvider::new("test-api-key".to_string()).unwrap();
// Test different WebSocket message types
let messages = vec![
Message::Text(r#"{"type": "status", "message": "Connected", "timestamp": "2024-01-15T09:30:00Z", "level": "info"}"#.to_string()),
Message::Binary(vec![0x01, 0x02, 0x03]),
Message::Ping(vec![0x01]),
Message::Pong(vec![0x01]),
Message::Close(None),
];
for message in messages {
let result = provider.handle_message(message).await;
assert!(result.is_ok());
}
}