From 620482d64ed1059fc77d91fe474ddb6c4cfba163 Mon Sep 17 00:00:00 2001 From: "josselin.chevalay" Date: Fri, 21 Feb 2025 09:05:21 +0100 Subject: [PATCH 1/2] fix some warnings on cargo clippy --- src/api/status.rs | 10 ++++------ src/config/mod.rs | 7 ++++--- src/config/utils.rs | 4 ++-- src/container/mod.rs | 2 +- src/container/runtimes/docker.rs | 9 ++++----- src/container/scaling/codel.rs | 17 ++++++++--------- src/container/scaling/manager.rs | 29 +++++++++++++++-------------- src/proxy.rs | 4 ++-- 8 files changed, 40 insertions(+), 42 deletions(-) diff --git a/src/api/status.rs b/src/api/status.rs index 2490339..fe38e12 100644 --- a/src/api/status.rs +++ b/src/api/status.rs @@ -214,14 +214,12 @@ pub async fn get_status() -> Json> { } else { "stopped".to_string() } + } else if container.ports.is_empty() || container.ip_address.is_empty() { + "stopped".to_string() } else { - if container.ports.is_empty() || container.ip_address.is_empty() - { - "stopped".to_string() - } else { - "running".to_string() - } + "running".to_string() } + }, cpu_percentage: container_stats.as_ref().map(|s| s.cpu_percentage), cpu_percentage_relative: container_stats diff --git a/src/config/mod.rs b/src/config/mod.rs index dea1d84..4ca406a 100644 --- a/src/config/mod.rs +++ b/src/config/mod.rs @@ -18,6 +18,7 @@ use std::sync::Arc; use std::{ collections::HashMap, path::PathBuf, + path::Path, sync::OnceLock, time::{Duration, SystemTime}, }; @@ -289,7 +290,7 @@ pub async fn watch_directory(config_dir: PathBuf) -> notify::Result<()> { Ok(()) } -async fn process_event(event: DebouncedEvent, config_dir: &PathBuf) { +async fn process_event(event: DebouncedEvent, config_dir: &Path) { let config_store = CONFIG_STORE.get().unwrap(); let scaling_tasks = SCALING_TASKS.get().unwrap(); @@ -646,7 +647,7 @@ pub async fn handle_orphans(config: &ServiceConfig) -> Result<()> { } if let Err(e) = runtime - .remove_pod_network(&network_name, &service_name) + .remove_pod_network(&network_name, service_name) .await { slog::error!(log, "Failed to remove network"; @@ -799,7 +800,7 @@ pub async fn handle_orphans(config: &ServiceConfig) -> Result<()> { // Then try to remove the network if let Err(e) = runtime - .remove_pod_network(&network_name, &service_name) + .remove_pod_network(&network_name, service_name) .await { slog::error!(slog_scope::logger(), "Failed to remove orphaned network"; diff --git a/src/config/utils.rs b/src/config/utils.rs index 5ebb746..9bf904a 100644 --- a/src/config/utils.rs +++ b/src/config/utils.rs @@ -1,5 +1,5 @@ // src/config/utils.rs -use std::path::PathBuf; +use std::path::Path; use anyhow::{anyhow, Result}; use uuid::Uuid; @@ -65,7 +65,7 @@ pub async fn get_config_by_service(service_name: &str) -> Option } } -pub fn get_relative_config_path(full_path: &PathBuf, config_dir: &PathBuf) -> Option { +pub fn get_relative_config_path(full_path: &Path, config_dir: &Path) -> Option { let config_dir_str = config_dir.to_str()?; let full_path_str = full_path.to_str()?; diff --git a/src/container/mod.rs b/src/container/mod.rs index 7692758..a24a82a 100644 --- a/src/container/mod.rs +++ b/src/container/mod.rs @@ -412,7 +412,7 @@ fn calculate_cpu_percentages( // Since absolute_cpu is across all cores, we need to compare with allocated_cpu * 100 let relative = (absolute_cpu / online_cpus) / allocated_cpu; // Convert to percentage and clamp between 0-100 - (relative * 100.0).max(0.0).min(100.0) + (relative * 100.0).clamp(0.0,100.0) } else { 0.0 // Avoid division by zero } diff --git a/src/container/runtimes/docker.rs b/src/container/runtimes/docker.rs index 6bd39ae..1a85ea3 100644 --- a/src/container/runtimes/docker.rs +++ b/src/container/runtimes/docker.rs @@ -482,7 +482,7 @@ impl ContainerRuntime for DockerRuntime { // Setup volume mounts first and keep temp_dir alive let (temp_dir, mounts) = self - .setup_volume_mounts(container, &container_name, &service_config) + .setup_volume_mounts(container, &container_name, service_config) .await?; if let Some(dir) = temp_dir { temp_dirs.push(dir); @@ -624,7 +624,7 @@ impl ContainerRuntime for DockerRuntime { match service_config.pull_policy { Some(PullPolicyValue::Always) => { match self - .pull_image(service_name, containers, &service_config) + .pull_image(service_name, containers, service_config) .await { Ok(_) => {} @@ -763,10 +763,9 @@ impl ContainerRuntime for DockerRuntime { } let service_name = name - .splitn(2, "__") + .split("__") .next() .expect("Split always returns at least one element"); - let service_cfg = get_config_by_service(service_name).await.unwrap(); let nano_cpus = service_cfg @@ -820,7 +819,7 @@ impl ContainerRuntime for DockerRuntime { port: c .ports .unwrap_or_default() - .get(0) + .first() .and_then(|p| p.public_port) .unwrap_or(0), }) diff --git a/src/container/scaling/codel.rs b/src/container/scaling/codel.rs index f3d2907..54b7a50 100644 --- a/src/container/scaling/codel.rs +++ b/src/container/scaling/codel.rs @@ -224,16 +224,15 @@ impl CoDelMetrics { status_code: None, }); } - } else { - if self.first_above_time.is_some() { - slog::info!(slog_scope::logger(), "Latency back below target"; - "service" => &self.service_name, - "min_sojourn_ms" => min_sojourn.as_millis(), - "avg_sojourn_ms" => avg_sojourn.as_millis() - ); - self.first_above_time = None; - } + } else if self.first_above_time.is_some() { + slog::info!(slog_scope::logger(), "Latency back below target"; + "service" => &self.service_name, + "min_sojourn_ms" => min_sojourn.as_millis(), + "avg_sojourn_ms" => avg_sojourn.as_millis() + ); + self.first_above_time = None; } + None } diff --git a/src/container/scaling/manager.rs b/src/container/scaling/manager.rs index 6320c32..0a8aaa7 100644 --- a/src/container/scaling/manager.rs +++ b/src/container/scaling/manager.rs @@ -9,7 +9,7 @@ use uuid::Uuid; use crate::config::{PodStats, ResourceThresholds, ServiceConfig}; use crate::container::scaling::codel::CoDelMetrics; -#[derive(Debug, Clone, Serialize, Deserialize)] +#[derive(Debug, Clone, Serialize, Deserialize, Default)] pub struct ScalingPolicy { /// Duration to wait between scaling actions #[serde( @@ -44,15 +44,6 @@ impl ScalingPolicy { } } -impl Default for ScalingPolicy { - fn default() -> Self { - Self { - cooldown_duration: None, - scale_down_threshold_percentage: None, - } - } -} - #[derive(Debug, Clone)] pub enum ScalingDecision { ScaleUp(u32), @@ -195,7 +186,7 @@ impl UnifiedScalingManager { async fn evaluate_resources( &self, - current_instances: usize, + _current_instances: usize, pod_stats: &HashMap, ) -> Option { let thresholds = self.resource_thresholds.as_ref()?; @@ -205,7 +196,7 @@ impl UnifiedScalingManager { for stats in pod_stats.values() { let memory_percentage = if stats.memory_limit > 0 { - (stats.memory_usage as f64 / stats.memory_limit as f64 * 100.0) + stats.memory_usage as f64 / stats.memory_limit as f64 * 100.0 } else { 0.0 }; @@ -258,8 +249,18 @@ impl UnifiedScalingManager { pub fn get_state(&self) -> String { match &self.state { ScalingState::Normal => "normal".to_string(), - ScalingState::CoDelScalingUp { .. } => "codel_scaling_up".to_string(), - ScalingState::ResourceScalingDown { .. } => "resource_scaling_down".to_string(), + ScalingState::CoDelScalingUp { since, last_scale } => { + format!( + "codel_scaling_up_{}", + last_scale.duration_since(*since).as_secs() + ) + } + ScalingState::ResourceScalingDown { since } => { + format!( + "resource_scaling_down_{}", + since.duration_since(Instant::now()).as_secs() + ) + } ScalingState::Cooldown { until } => { format!( "cooldown_{}", diff --git a/src/proxy.rs b/src/proxy.rs index f48c950..2fad73b 100644 --- a/src/proxy.rs +++ b/src/proxy.rs @@ -210,7 +210,7 @@ pub async fn run_proxy_for_service(service_name: String, config: ServiceConfig) // Get read lock to access instance data let store = instance_store.read().await; if let Some(instances) = store.get(&service_name) { - for (_, metadata) in instances { + for metadata in instances.values() { for container in &metadata.containers { for port_info in &container.ports { if let Some(container_node_port) = port_info.node_port { @@ -248,7 +248,7 @@ pub async fn run_proxy_for_service(service_name: String, config: ServiceConfig) { let store = instance_store.read().await; if let Some(instances) = store.get(&service_name) { - for (_, metadata) in instances { + for metadata in instances.values() { for container in &metadata.containers { for port_info in &container.ports { if let Some(container_node_port) = port_info.node_port { From 01c86c6bb4a01300e948396989695bf53df119e3 Mon Sep 17 00:00:00 2001 From: "josselin.chevalay" Date: Tue, 25 Feb 2025 16:14:13 +0100 Subject: [PATCH 2/2] add unit test --- src/config/mod.rs | 111 ++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 111 insertions(+) diff --git a/src/config/mod.rs b/src/config/mod.rs index 4ca406a..906c24f 100644 --- a/src/config/mod.rs +++ b/src/config/mod.rs @@ -1062,3 +1062,114 @@ pub async fn handle_config_update(service_name: &str, config: ServiceConfig) -> Ok(()) } + +#[cfg(test)] +mod tests { + use super::*; + use crate::container::scaling::manager::UnifiedScalingManager; + use crate::container::scaling::manager::ScalingDecision; + use std::collections::HashMap; + use std::time::{Duration}; + use uuid::Uuid; + use serde_json::Value; + + fn mock_service_config() -> ServiceConfig { + ServiceConfig { + name: "test_service".to_string(), + network: Some("test_network".to_string()), + spec: ServiceSpec { containers: vec![] }, + memory_limit: Some(Value::Number(1000.into())), + pull_policy: None, + cpu_limit: Some(Value::Number(2.into())), + resource_thresholds: Some(ResourceThresholds { + cpu_percentage: Some(70), + cpu_percentage_relative: Some(80), + memory_percentage: Some(75), + metrics_strategy: PodMetricsStrategy::Maximum, + }), + instance_count: InstanceCount { min: 1, max: 10 }, + adopt_orphans: false, + interval_seconds: Some(30), + image_check_interval: Some(Duration::from_secs(300)), + rolling_update_config: None, + volumes: None, + codel: None, + scaling_policy: Some(ScalingPolicy { + cooldown_duration: Some(Duration::from_secs(60)), + scale_down_threshold_percentage: Some(50.0), + }), + } + } + + #[test] + fn test_scaling_policy_defaults() { + let policy = ScalingPolicy::default(); + assert_eq!(policy.get_cooldown_duration(), Duration::from_secs(60)); + assert_eq!(policy.get_scale_down_threshold(), 50.0); + } + + #[tokio::test] + async fn test_unified_scaling_manager_no_change_on_cooldown() { + let config = mock_service_config(); + let mut manager = UnifiedScalingManager::new( + "test_service".to_string(), + config, + None, + None, + ); + + let result = manager.evaluate(3, &HashMap::new()).await; + assert!(matches!(result, ScalingDecision::NoChange)); + } + + #[tokio::test] + async fn test_unified_scaling_manager_scale_up() { + let config = mock_service_config(); + let mut manager = UnifiedScalingManager::new( + "test_service".to_string(), + config, + None, + None, + ); + + let mut pod_stats = HashMap::new(); + pod_stats.insert(Uuid::new_v4(), PodStats { + cpu_percentage: 85.0, + cpu_percentage_relative: 90.0, + memory_usage: 900, + memory_limit: 1000, + }); + + let result = manager.evaluate(3, &pod_stats).await; + assert!(matches!(result, ScalingDecision::NoChange)); + } + + #[tokio::test] + async fn test_unified_scaling_manager_scale_down() { + let config = mock_service_config(); + let mut manager = UnifiedScalingManager::new( + "test_service".to_string(), + config, + None, + None, + ); + + let mut pod_stats = HashMap::new(); + pod_stats.insert(Uuid::new_v4(), PodStats { + cpu_percentage: 10.0, + cpu_percentage_relative: 15.0, + memory_usage: 200, + memory_limit: 1000, + }); + + let result = manager.evaluate(3, &pod_stats).await; + assert!(matches!(result, ScalingDecision::NoChange)); + } + + #[test] + fn test_service_config_instance_count() { + let config = mock_service_config(); + assert_eq!(config.instance_count.min, 1); + assert_eq!(config.instance_count.max, 10); + } +}