//! Service Management for E2E Testing //! //! Provides orchestration capabilities for starting, stopping, and managing //! all services required for end-to-end testing of the Foxhunt system. use anyhow::{Context, Result}; use std::collections::HashMap; use std::path::PathBuf; use std::process::{Child, Command, Stdio}; use std::time::Duration; use tokio::time::sleep; use tracing::{debug, error, info, warn}; /// Service type enumeration for orchestrator #[derive(Debug, Clone, PartialEq, Eq, Hash)] pub enum ServiceType { ApiGateway, TradingService, BacktestingService, MLTrainingService, Database, } impl ServiceType { pub fn as_str(&self) -> &str { match self { ServiceType::ApiGateway => "api", ServiceType::TradingService => "trading", ServiceType::BacktestingService => "backtesting", ServiceType::MLTrainingService => "ml_training", ServiceType::Database => "database", } } } /// Service configuration for orchestrator #[derive(Debug, Clone)] pub struct ServiceConfig { pub service_type: ServiceType, pub executable_path: String, pub port: u16, pub health_endpoint: String, pub startup_timeout: Duration, pub environment: HashMap, pub working_directory: PathBuf, pub log_file: Option, } /// Legacy service configuration (kept for backward compatibility) #[derive(Debug, Clone)] pub struct LegacyServiceConfig { pub name: String, pub binary_name: String, pub port: u16, pub args: Vec, pub env_vars: HashMap, pub startup_timeout: Duration, } /// Service manager for orchestrating all Foxhunt services #[derive(Debug)] pub struct ServiceManager { services: HashMap, processes: HashMap, base_path: std::path::PathBuf, } impl Default for ServiceManager { fn default() -> Self { Self::new() } } impl ServiceManager { /// Create a new service manager pub fn new() -> Self { Self { services: HashMap::new(), processes: HashMap::new(), base_path: std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")), } } /// Create a new service manager with legacy configs pub fn new_with_configs() -> Result { let base_path = std::env::current_dir().context("Failed to get current directory")?; let services = Self::create_service_configs(); Ok(Self { services, processes: HashMap::new(), base_path, }) } /// Start a service using orchestrator ServiceConfig pub async fn start_service(&mut self, config: ServiceConfig) -> Result<()> { tracing::info!("Starting service: {:?}", config.service_type); // For now, just track that we tried to start it // Full implementation would spawn the actual service process Ok(()) } /// Start all services pub async fn start_all_services(&mut self) -> Result<()> { info!("🚀 Starting all services for E2E testing"); // Start services in dependency order let service_order = vec!["trading_service", "backtesting_service"]; for service_name in service_order { if let Some(config) = self.services.get(service_name) { self.start_legacy_service(config.clone()) .await .with_context(|| format!("Failed to start service: {}", service_name))?; } else { warn!("Service configuration not found: {}", service_name); } } info!("✅ All services started successfully"); Ok(()) } /// Stop all services pub async fn stop_all_services(&mut self) -> Result<()> { info!("🛑 Stopping all services"); // Stop in reverse order let service_order = vec!["backtesting_service", "trading_service"]; for service_name in service_order { if let Some(mut process) = self.processes.remove(service_name) { info!("Stopping {}", service_name); match process.kill() { Ok(_) => { info!("✅ {} stopped", service_name); }, Err(e) => { warn!("Failed to kill {}: {}", service_name, e); }, } // Wait for process to exit match process.wait() { Ok(_) => debug!("{} process exited", service_name), Err(e) => warn!("Error waiting for {} to exit: {}", service_name, e), } } } info!("✅ All services stopped"); Ok(()) } /// Start a specific service (legacy) async fn start_legacy_service(&mut self, config: LegacyServiceConfig) -> Result<()> { info!("🔧 Starting service: {}", config.name); // Check if binary exists let binary_path = self .base_path .join("target/release") .join(&config.binary_name); if !binary_path.exists() { // Try debug build let debug_binary_path = self .base_path .join("target/debug") .join(&config.binary_name); if !debug_binary_path.exists() { return Err(anyhow::anyhow!( "Service binary not found: {} (tried both release and debug)", config.binary_name )); } } // Build command let mut command = Command::new("cargo"); command .arg("run") .arg("--bin") .arg(&config.binary_name) .current_dir(&self.base_path) .stdout(Stdio::piped()) .stderr(Stdio::piped()); // Add arguments for arg in &config.args { command.arg(arg); } // Set environment variables for (key, value) in &config.env_vars { command.env(key, value); } // Set common environment variables command .env("RUST_LOG", "info") .env("FOXHUNT_TEST_MODE", "true") .env( "DATABASE_URL", std::env::var("DATABASE_URL") .unwrap_or_else(|_| "postgresql://localhost/foxhunt_test".to_string()), ); debug!("Starting command: {:?}", command); // Start the process let child = command .spawn() .with_context(|| format!("Failed to spawn {}", config.name))?; self.processes.insert(config.name.clone(), child); // Wait for service to be ready info!( "⏳ Waiting for {} to be ready on port {}...", config.name, config.port ); self.wait_for_legacy_service_ready(&config) .await .with_context(|| format!("Service {} failed to become ready", config.name))?; info!("✅ Service {} started successfully", config.name); Ok(()) } /// Wait for a service to be ready async fn wait_for_legacy_service_ready(&self, config: &LegacyServiceConfig) -> Result<()> { use tokio::net::TcpStream; let timeout = config.startup_timeout; let check_interval = Duration::from_millis(500); let start_time = std::time::Instant::now(); while start_time.elapsed() < timeout { match TcpStream::connect(("127.0.0.1", config.port)).await { Ok(_) => { debug!("Service {} is ready on port {}", config.name, config.port); return Ok(()); }, Err(_) => { debug!("Service {} not ready yet, retrying...", config.name); sleep(check_interval).await; }, } } Err(anyhow::anyhow!( "Service {} failed to become ready within {:?}", config.name, timeout )) } /// Create service configurations fn create_service_configs() -> HashMap { let mut services = HashMap::new(); // Trading Service services.insert( "trading_service".to_string(), LegacyServiceConfig { name: "trading_service".to_string(), binary_name: "trading_service".to_string(), port: 50051, args: vec![], env_vars: HashMap::new(), startup_timeout: Duration::from_secs(30), }, ); // Backtesting Service services.insert( "backtesting_service".to_string(), LegacyServiceConfig { name: "backtesting_service".to_string(), binary_name: "backtesting_service".to_string(), port: 50052, args: vec![], env_vars: HashMap::new(), startup_timeout: Duration::from_secs(30), }, ); services } /// Check if a service is running pub fn is_service_running(&self, service_name: &str) -> bool { self.processes.contains_key(service_name) } /// Get service status pub async fn get_service_status(&self, service_name: &str) -> ServiceStatus { if let Some(config) = self.services.get(service_name) { // Check if process is running let process_running = self.is_service_running(service_name); // Check if port is accessible let port_accessible = (tokio::net::TcpStream::connect(("127.0.0.1", config.port)).await).is_ok(); if process_running && port_accessible { ServiceStatus::Running } else if process_running { ServiceStatus::Starting } else { ServiceStatus::Stopped } } else { ServiceStatus::NotFound } } /// Get all service statuses pub async fn get_all_service_statuses(&self) -> HashMap { let mut statuses = HashMap::new(); for service_name in self.services.keys() { let status = self.get_service_status(service_name).await; statuses.insert(service_name.clone(), status); } statuses } } /// Service status enumeration #[derive(Debug, Clone, PartialEq)] pub enum ServiceStatus { Running, Starting, Stopped, NotFound, } impl std::fmt::Display for ServiceStatus { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { match self { ServiceStatus::Running => write!(f, "Running"), ServiceStatus::Starting => write!(f, "Starting"), ServiceStatus::Stopped => write!(f, "Stopped"), ServiceStatus::NotFound => write!(f, "Not Found"), } } } /// Cleanup implementation for ServiceManager impl Drop for ServiceManager { fn drop(&mut self) { // Try to stop all services on drop for (service_name, mut process) in self.processes.drain() { info!("Cleaning up service: {}", service_name); if let Err(e) = process.kill() { error!("Failed to kill service {}: {}", service_name, e); } } } } /// Helper function to check if all required binaries exist pub fn check_service_binaries() -> Result> { let base_path = std::env::current_dir().context("Failed to get current directory")?; let required_binaries = vec!["trading_service", "backtesting_service"]; let mut missing_binaries = Vec::new(); for binary in &required_binaries { let release_path = base_path.join("target/release").join(binary); let debug_path = base_path.join("target/debug").join(binary); if !release_path.exists() && !debug_path.exists() { missing_binaries.push(binary.to_string()); } } Ok(missing_binaries) } /// Build all required service binaries pub async fn build_service_binaries() -> Result<()> { info!("🔨 Building service binaries for E2E testing"); let binaries = vec!["trading_service", "backtesting_service"]; for binary in &binaries { info!("Building {}", binary); let output = tokio::process::Command::new("cargo") .arg("build") .arg("--bin") .arg(binary) .arg("--release") .output() .await .with_context(|| format!("Failed to build {}", binary))?; if !output.status.success() { let stderr = String::from_utf8_lossy(&output.stderr); return Err(anyhow::anyhow!("Failed to build {}: {}", binary, stderr)); } info!("✅ Built {}", binary); } info!("✅ All service binaries built successfully"); Ok(()) } #[cfg(test)] mod tests { use super::*; #[test] fn test_service_config_creation() { let services = ServiceManager::create_service_configs(); assert!(services.contains_key("trading_service")); assert!(services.contains_key("backtesting_service")); let trading_config = services.get("trading_service").unwrap(); assert_eq!(trading_config.port, 50051); assert_eq!(trading_config.binary_name, "trading_service"); } #[test] fn test_service_status_display() { assert_eq!(ServiceStatus::Running.to_string(), "Running"); assert_eq!(ServiceStatus::Starting.to_string(), "Starting"); assert_eq!(ServiceStatus::Stopped.to_string(), "Stopped"); assert_eq!(ServiceStatus::NotFound.to_string(), "Not Found"); } #[tokio::test] async fn test_service_manager_creation() { // ServiceManager::new() returns Self directly, not Result let manager = ServiceManager::new(); // Check that the manager has been initialized (empty services is expected for new()) assert_eq!(manager.services.len(), 0); // Test with configs let manager_with_configs = ServiceManager::new_with_configs(); assert!(manager_with_configs.is_ok()); let manager = manager_with_configs.unwrap(); assert!(manager.services.len() >= 2); } }