diff --git a/Cargo.lock b/Cargo.lock index 750df5563..fdedbe964 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -8120,7 +8120,7 @@ checksum = "9774ba4a74de5f7b1c1451ed6cd5285a32eddb5cccb8cc655a4e50009e06477f" [[package]] name = "s3s" version = "0.13.0-alpha.3" -source = "git+https://github.com/s3s-project/s3s.git?rev=218000387f4c3e67ad478bfc4587931f88b37006#218000387f4c3e67ad478bfc4587931f88b37006" +source = "git+https://github.com/s3s-project/s3s.git?rev=d3f8f8fcb1a1e0af9345498ae2b6eea6a18ef152#d3f8f8fcb1a1e0af9345498ae2b6eea6a18ef152" dependencies = [ "arc-swap", "arrayvec", diff --git a/Cargo.toml b/Cargo.toml index 237235c64..c5fb5a305 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -236,7 +236,7 @@ rumqttc = { version = "0.25.1" } rustix = { version = "1.1.4", features = ["fs"] } rust-embed = { version = "8.11.0" } rustc-hash = { version = "2.1.1" } -s3s = { version = "0.13.0-alpha.3", features = ["minio"], git = "https://github.com/s3s-project/s3s.git", rev = "218000387f4c3e67ad478bfc4587931f88b37006" } +s3s = { version = "0.13.0-alpha.3", features = ["minio"], git = "https://github.com/s3s-project/s3s.git", rev = "d3f8f8fcb1a1e0af9345498ae2b6eea6a18ef152" } serial_test = "3.4.0" shadow-rs = { version = "1.7.0", default-features = false } siphasher = "1.0.2" diff --git a/crates/config/src/constants/app.rs b/crates/config/src/constants/app.rs index 7b1aa2756..a39b71ede 100644 --- a/crates/config/src/constants/app.rs +++ b/crates/config/src/constants/app.rs @@ -150,8 +150,8 @@ pub const DEFAULT_CONSOLE_ADDRESS: &str = concat!(":", DEFAULT_CONSOLE_PORT); /// Default region for rustfs /// This is the default region for rustfs. /// It is used to identify the region of the application. -/// Default value: rustfs-global-0 -pub const RUSTFS_REGION: &str = "rustfs-global-0"; +/// Default value: us-east-1 +pub const RUSTFS_REGION: &str = "us-east-1"; /// Default log filename for rustfs /// This is the default log filename for rustfs. diff --git a/rustfs/src/config/mod.rs b/rustfs/src/config/mod.rs index 9a3f1e84e..c68345161 100644 --- a/rustfs/src/config/mod.rs +++ b/rustfs/src/config/mod.rs @@ -284,7 +284,7 @@ impl Config { .trim() .to_string(); - // Region is optional, but if not set, we should default to "rustfs-global-0" for signing compatibility with AWS S3 clients + // Region is optional, but if not set, we should default to "us-east-1" for signing compatibility with AWS S3 clients let region = region.or_else(|| Some(RUSTFS_REGION.to_string())); Ok(Config { diff --git a/rustfs/src/init.rs b/rustfs/src/init.rs index 09f8d77d5..add0814d3 100644 --- a/rustfs/src/init.rs +++ b/rustfs/src/init.rs @@ -14,7 +14,7 @@ use crate::storage::{process_lambda_configurations, process_queue_configurations, process_topic_configurations}; use crate::{admin, config, version}; -use rustfs_config::{DEFAULT_UPDATE_CHECK, ENV_UPDATE_CHECK}; +use rustfs_config::{DEFAULT_UPDATE_CHECK, ENV_UPDATE_CHECK, RUSTFS_REGION}; use rustfs_ecstore::bucket::metadata_sys; use rustfs_notify::notifier_global; use rustfs_targets::arn::{ARN, TargetIDError}; @@ -85,6 +85,14 @@ pub(crate) fn init_update_check() { }); } +/// Helper function to parse ARN string to target ID +/// Converts an ARN string to a target ID, or returns an error if parsing fails +fn arn_to_target_id(arn_str: &str) -> Result { + ARN::parse(arn_str) + .map(|arn| arn.target_id) + .map_err(|e| TargetIDError::InvalidFormat(e.to_string())) +} + /// Add existing bucket notification configurations to the global notifier system. /// This function retrieves notification configurations for each bucket /// and registers the corresponding event rules with the notifier system. @@ -99,8 +107,11 @@ pub(crate) async fn add_bucket_notification_configuration(buckets: Vec) .filter(|r| !r.as_str().is_empty()) .map(|r| r.as_str()) .unwrap_or_else(|| { - warn!("Global region is not set; attempting notification configuration for all buckets with an empty region."); - "" + warn!( + "Global region is not set; attempting notification configuration for all buckets using default region '{}'.", + RUSTFS_REGION + ); + RUSTFS_REGION }); for bucket in buckets.iter() { let has_notification_config = metadata_sys::get_notification_config(bucket).await.unwrap_or_else(|err| { @@ -116,26 +127,16 @@ pub(crate) async fn add_bucket_notification_configuration(buckets: Vec) "Bucket '{}' has existing notification configuration: {:?}", bucket, cfg); let mut event_rules = Vec::new(); - if let Err(e) = process_queue_configurations(&mut event_rules, cfg.queue_configurations.clone(), |arn_str| { - ARN::parse(arn_str) - .map(|arn| arn.target_id) - .map_err(|e| TargetIDError::InvalidFormat(e.to_string())) - }) { + if let Err(e) = process_queue_configurations(&mut event_rules, cfg.queue_configurations.clone(), arn_to_target_id) + { error!("Failed to parse queue notification config for bucket '{}': {:?}", bucket, e); } - if let Err(e) = process_topic_configurations(&mut event_rules, cfg.topic_configurations.clone(), |arn_str| { - ARN::parse(arn_str) - .map(|arn| arn.target_id) - .map_err(|e| TargetIDError::InvalidFormat(e.to_string())) - }) { + if let Err(e) = process_topic_configurations(&mut event_rules, cfg.topic_configurations.clone(), arn_to_target_id) + { error!("Failed to parse topic notification config for bucket '{}': {:?}", bucket, e); } if let Err(e) = - process_lambda_configurations(&mut event_rules, cfg.lambda_function_configurations.clone(), |arn_str| { - ARN::parse(arn_str) - .map(|arn| arn.target_id) - .map_err(|e| TargetIDError::InvalidFormat(e.to_string())) - }) + process_lambda_configurations(&mut event_rules, cfg.lambda_function_configurations.clone(), arn_to_target_id) { error!("Failed to parse lambda notification config for bucket '{}': {:?}", bucket, e); } @@ -157,6 +158,80 @@ pub(crate) async fn add_bucket_notification_configuration(buckets: Vec) } } +/// Build KMS configuration for local backend +fn build_local_kms_config(cfg: &config::Config) -> std::io::Result { + let key_dir = cfg + .kms_key_dir + .as_ref() + .ok_or_else(|| Error::other("KMS key directory is required for local backend"))?; + + Ok(rustfs_kms::config::KmsConfig { + backend: rustfs_kms::config::KmsBackend::Local, + backend_config: rustfs_kms::config::BackendConfig::Local(rustfs_kms::config::LocalConfig { + key_dir: std::path::PathBuf::from(key_dir), + master_key: None, + file_permissions: Some(0o600), + }), + default_key_id: cfg.kms_default_key_id.clone(), + timeout: std::time::Duration::from_secs(30), + retry_attempts: 3, + enable_cache: true, + cache_config: rustfs_kms::config::CacheConfig::default(), + }) +} + +/// Build KMS configuration for Vault backend +fn build_vault_kms_config(cfg: &config::Config) -> std::io::Result { + let vault_address = cfg + .kms_vault_address + .as_ref() + .ok_or_else(|| Error::other("Vault address is required for vault backend"))?; + let vault_token = cfg + .kms_vault_token + .as_ref() + .ok_or_else(|| Error::other("Vault token is required for vault backend"))?; + + Ok(rustfs_kms::config::KmsConfig { + backend: rustfs_kms::config::KmsBackend::Vault, + backend_config: rustfs_kms::config::BackendConfig::Vault(Box::new(rustfs_kms::config::VaultConfig { + address: vault_address.clone(), + auth_method: rustfs_kms::config::VaultAuthMethod::Token { + token: vault_token.clone(), + }, + namespace: None, + mount_path: "transit".to_string(), + kv_mount: "secret".to_string(), + key_path_prefix: "rustfs/kms/keys".to_string(), + tls: None, + })), + default_key_id: cfg.kms_default_key_id.clone(), + timeout: std::time::Duration::from_secs(30), + retry_attempts: 3, + enable_cache: true, + cache_config: rustfs_kms::config::CacheConfig::default(), + }) +} + +/// Configure and start KMS service +async fn configure_and_start_kms( + service_manager: &std::sync::Arc, + kms_config: rustfs_kms::config::KmsConfig, + config_source: &str, +) -> std::io::Result<()> { + service_manager + .configure(kms_config) + .await + .map_err(|e| Error::other(format!("Failed to configure KMS: {e}")))?; + + service_manager + .start() + .await + .map_err(|e| Error::other(format!("Failed to start KMS: {e}")))?; + + info!("KMS service configured and started successfully from {}", config_source); + Ok(()) +} + /// Initialize KMS system and configure if enabled /// /// This function initializes the global KMS service manager. If KMS is enabled @@ -178,72 +253,12 @@ pub(crate) async fn init_kms_system(config: &config::Config) -> std::io::Result< // Create KMS configuration from command line options let kms_config = match config.kms_backend.as_str() { - "local" => { - let key_dir = config - .kms_key_dir - .as_ref() - .ok_or_else(|| Error::other("KMS key directory is required for local backend"))?; - - rustfs_kms::config::KmsConfig { - backend: rustfs_kms::config::KmsBackend::Local, - backend_config: rustfs_kms::config::BackendConfig::Local(rustfs_kms::config::LocalConfig { - key_dir: std::path::PathBuf::from(key_dir), - master_key: None, - file_permissions: Some(0o600), - }), - default_key_id: config.kms_default_key_id.clone(), - timeout: std::time::Duration::from_secs(30), - retry_attempts: 3, - enable_cache: true, - cache_config: rustfs_kms::config::CacheConfig::default(), - } - } - "vault" => { - let vault_address = config - .kms_vault_address - .as_ref() - .ok_or_else(|| Error::other("Vault address is required for vault backend"))?; - let vault_token = config - .kms_vault_token - .as_ref() - .ok_or_else(|| Error::other("Vault token is required for vault backend"))?; - - rustfs_kms::config::KmsConfig { - backend: rustfs_kms::config::KmsBackend::Vault, - backend_config: rustfs_kms::config::BackendConfig::Vault(Box::new(rustfs_kms::config::VaultConfig { - address: vault_address.clone(), - auth_method: rustfs_kms::config::VaultAuthMethod::Token { - token: vault_token.clone(), - }, - namespace: None, - mount_path: "transit".to_string(), - kv_mount: "secret".to_string(), - key_path_prefix: "rustfs/kms/keys".to_string(), - tls: None, - })), - default_key_id: config.kms_default_key_id.clone(), - timeout: std::time::Duration::from_secs(30), - retry_attempts: 3, - enable_cache: true, - cache_config: rustfs_kms::config::CacheConfig::default(), - } - } + "local" => build_local_kms_config(config)?, + "vault" => build_vault_kms_config(config)?, _ => return Err(Error::other(format!("Unsupported KMS backend: {}", config.kms_backend))), }; - // Configure the KMS service - service_manager - .configure(kms_config) - .await - .map_err(|e| Error::other(format!("Failed to configure KMS: {e}")))?; - - // Start the KMS service - service_manager - .start() - .await - .map_err(|e| Error::other(format!("Failed to start KMS: {e}")))?; - - info!("KMS service configured and started successfully from command line options"); + configure_and_start_kms(&service_manager, kms_config, "command line options").await?; } else { // Try to load persisted KMS configuration from cluster storage info!("Attempting to load persisted KMS configuration from cluster storage..."); @@ -252,17 +267,9 @@ pub(crate) async fn init_kms_system(config: &config::Config) -> std::io::Result< info!("Found persisted KMS configuration, attempting to configure and start service..."); // Configure the KMS service with persisted config - match service_manager.configure(persisted_config).await { + match configure_and_start_kms(&service_manager, persisted_config, "persisted configuration").await { Ok(()) => { - // Start the KMS service - match service_manager.start().await { - Ok(()) => { - info!("KMS service configured and started successfully from persisted configuration"); - } - Err(e) => { - warn!("Failed to start KMS with persisted configuration: {}", e); - } - } + info!("KMS service configured and started successfully from persisted configuration"); } Err(e) => { warn!("Failed to configure KMS with persisted configuration: {}", e); @@ -318,6 +325,46 @@ pub(crate) fn init_buffer_profile_system(config: &config::Config) { } } +/// Parse and normalize server address for FTP/FTPS +/// Forces IPv4 binding to avoid libunftp IPv6 compatibility issues +#[allow(dead_code)] +async fn parse_and_normalize_server_address( + address_str: &str, +) -> Result> { + let addr = rustfs_utils::net::parse_and_resolve_address(address_str) + .map_err(|e| format!("Invalid server address '{address_str}': {e}"))?; + + // Force IPv4 binding to avoid libunftp IPv6 compatibility issues + let normalized_addr = if addr.is_ipv6() { + std::net::SocketAddr::new(std::net::IpAddr::V4(std::net::Ipv4Addr::UNSPECIFIED), addr.port()) + } else { + addr + }; + + Ok(normalized_addr) +} + +/// Start FTP/FTPS server in background with shutdown support +/// # Arguments +/// * `server` - The FTP/FTPS server instance +/// * `protocol_name` - Name of the protocol (e.g., "FTP", "FTPS") +#[allow(dead_code)] +fn spawn_server(server: S, protocol_name: &'static str) -> tokio::sync::broadcast::Sender<()> +where + S: std::future::Future>> + Send + 'static, +{ + let (shutdown_tx, _) = tokio::sync::broadcast::channel(1); + + tokio::spawn(async move { + if let Err(e) = server.await { + error!("{} server error: {}", protocol_name, e); + } + info!("{} server shutdown completed", protocol_name); + }); + + shutdown_tx +} + /// Initialize the FTP system /// /// This function initializes the FTP server (non-encrypted) if enabled in the configuration. @@ -338,14 +385,7 @@ pub async fn init_ftp_system() -> Result Result Result>, Box> { @@ -418,14 +458,7 @@ pub async fn init_ftps_system() -> Result