diff --git a/crates/config/src/audit/kafka.rs b/crates/config/src/audit/kafka.rs index b67d566a2..5604690b2 100644 --- a/crates/config/src/audit/kafka.rs +++ b/crates/config/src/audit/kafka.rs @@ -21,10 +21,14 @@ pub const ENV_AUDIT_KAFKA_TLS_ENABLE: &str = "RUSTFS_AUDIT_KAFKA_TLS_ENABLE"; pub const ENV_AUDIT_KAFKA_TLS_CA: &str = "RUSTFS_AUDIT_KAFKA_TLS_CA"; pub const ENV_AUDIT_KAFKA_TLS_CLIENT_CERT: &str = "RUSTFS_AUDIT_KAFKA_TLS_CLIENT_CERT"; pub const ENV_AUDIT_KAFKA_TLS_CLIENT_KEY: &str = "RUSTFS_AUDIT_KAFKA_TLS_CLIENT_KEY"; +pub const ENV_AUDIT_KAFKA_SASL_ENABLE: &str = "RUSTFS_AUDIT_KAFKA_SASL_ENABLE"; +pub const ENV_AUDIT_KAFKA_SASL_MECHANISM: &str = "RUSTFS_AUDIT_KAFKA_SASL_MECHANISM"; +pub const ENV_AUDIT_KAFKA_SASL_USERNAME: &str = "RUSTFS_AUDIT_KAFKA_SASL_USERNAME"; +pub const ENV_AUDIT_KAFKA_SASL_PASSWORD: &str = "RUSTFS_AUDIT_KAFKA_SASL_PASSWORD"; pub const ENV_AUDIT_KAFKA_QUEUE_DIR: &str = "RUSTFS_AUDIT_KAFKA_QUEUE_DIR"; pub const ENV_AUDIT_KAFKA_QUEUE_LIMIT: &str = "RUSTFS_AUDIT_KAFKA_QUEUE_LIMIT"; -pub const ENV_AUDIT_KAFKA_KEYS: &[&str; 10] = &[ +pub const ENV_AUDIT_KAFKA_KEYS: &[&str; 14] = &[ ENV_AUDIT_KAFKA_ENABLE, ENV_AUDIT_KAFKA_BROKERS, ENV_AUDIT_KAFKA_TOPIC, @@ -33,6 +37,10 @@ pub const ENV_AUDIT_KAFKA_KEYS: &[&str; 10] = &[ ENV_AUDIT_KAFKA_TLS_CA, ENV_AUDIT_KAFKA_TLS_CLIENT_CERT, ENV_AUDIT_KAFKA_TLS_CLIENT_KEY, + ENV_AUDIT_KAFKA_SASL_ENABLE, + ENV_AUDIT_KAFKA_SASL_MECHANISM, + ENV_AUDIT_KAFKA_SASL_USERNAME, + ENV_AUDIT_KAFKA_SASL_PASSWORD, ENV_AUDIT_KAFKA_QUEUE_DIR, ENV_AUDIT_KAFKA_QUEUE_LIMIT, ]; @@ -47,6 +55,10 @@ pub const AUDIT_KAFKA_KEYS: &[&str] = &[ crate::KAFKA_TLS_CA, crate::KAFKA_TLS_CLIENT_CERT, crate::KAFKA_TLS_CLIENT_KEY, + crate::KAFKA_SASL_ENABLE, + crate::KAFKA_SASL_MECHANISM, + crate::KAFKA_SASL_USERNAME, + crate::KAFKA_SASL_PASSWORD, crate::KAFKA_QUEUE_DIR, crate::KAFKA_QUEUE_LIMIT, crate::COMMENT_KEY, diff --git a/crates/config/src/constants/targets.rs b/crates/config/src/constants/targets.rs index 40344589d..0d3e14a32 100644 --- a/crates/config/src/constants/targets.rs +++ b/crates/config/src/constants/targets.rs @@ -49,6 +49,10 @@ pub const KAFKA_TLS_ENABLE: &str = "tls_enable"; pub const KAFKA_TLS_CA: &str = "tls_ca"; pub const KAFKA_TLS_CLIENT_CERT: &str = "tls_client_cert"; pub const KAFKA_TLS_CLIENT_KEY: &str = "tls_client_key"; +pub const KAFKA_SASL_ENABLE: &str = "sasl_enable"; +pub const KAFKA_SASL_MECHANISM: &str = "sasl_mechanism"; +pub const KAFKA_SASL_USERNAME: &str = "sasl_username"; +pub const KAFKA_SASL_PASSWORD: &str = "sasl_password"; pub const AMQP_URL: &str = "url"; pub const AMQP_EXCHANGE: &str = "exchange"; diff --git a/crates/config/src/notify/kafka.rs b/crates/config/src/notify/kafka.rs index e8112096e..13c18c76f 100644 --- a/crates/config/src/notify/kafka.rs +++ b/crates/config/src/notify/kafka.rs @@ -22,6 +22,10 @@ pub const NOTIFY_KAFKA_KEYS: &[&str] = &[ crate::KAFKA_TLS_CA, crate::KAFKA_TLS_CLIENT_CERT, crate::KAFKA_TLS_CLIENT_KEY, + crate::KAFKA_SASL_ENABLE, + crate::KAFKA_SASL_MECHANISM, + crate::KAFKA_SASL_USERNAME, + crate::KAFKA_SASL_PASSWORD, crate::KAFKA_QUEUE_DIR, crate::KAFKA_QUEUE_LIMIT, crate::COMMENT_KEY, @@ -36,10 +40,14 @@ pub const ENV_NOTIFY_KAFKA_TLS_ENABLE: &str = "RUSTFS_NOTIFY_KAFKA_TLS_ENABLE"; pub const ENV_NOTIFY_KAFKA_TLS_CA: &str = "RUSTFS_NOTIFY_KAFKA_TLS_CA"; pub const ENV_NOTIFY_KAFKA_TLS_CLIENT_CERT: &str = "RUSTFS_NOTIFY_KAFKA_TLS_CLIENT_CERT"; pub const ENV_NOTIFY_KAFKA_TLS_CLIENT_KEY: &str = "RUSTFS_NOTIFY_KAFKA_TLS_CLIENT_KEY"; +pub const ENV_NOTIFY_KAFKA_SASL_ENABLE: &str = "RUSTFS_NOTIFY_KAFKA_SASL_ENABLE"; +pub const ENV_NOTIFY_KAFKA_SASL_MECHANISM: &str = "RUSTFS_NOTIFY_KAFKA_SASL_MECHANISM"; +pub const ENV_NOTIFY_KAFKA_SASL_USERNAME: &str = "RUSTFS_NOTIFY_KAFKA_SASL_USERNAME"; +pub const ENV_NOTIFY_KAFKA_SASL_PASSWORD: &str = "RUSTFS_NOTIFY_KAFKA_SASL_PASSWORD"; pub const ENV_NOTIFY_KAFKA_QUEUE_DIR: &str = "RUSTFS_NOTIFY_KAFKA_QUEUE_DIR"; pub const ENV_NOTIFY_KAFKA_QUEUE_LIMIT: &str = "RUSTFS_NOTIFY_KAFKA_QUEUE_LIMIT"; -pub const ENV_NOTIFY_KAFKA_KEYS: &[&str; 10] = &[ +pub const ENV_NOTIFY_KAFKA_KEYS: &[&str; 14] = &[ ENV_NOTIFY_KAFKA_ENABLE, ENV_NOTIFY_KAFKA_BROKERS, ENV_NOTIFY_KAFKA_TOPIC, @@ -48,6 +56,10 @@ pub const ENV_NOTIFY_KAFKA_KEYS: &[&str; 10] = &[ ENV_NOTIFY_KAFKA_TLS_CA, ENV_NOTIFY_KAFKA_TLS_CLIENT_CERT, ENV_NOTIFY_KAFKA_TLS_CLIENT_KEY, + ENV_NOTIFY_KAFKA_SASL_ENABLE, + ENV_NOTIFY_KAFKA_SASL_MECHANISM, + ENV_NOTIFY_KAFKA_SASL_USERNAME, + ENV_NOTIFY_KAFKA_SASL_PASSWORD, ENV_NOTIFY_KAFKA_QUEUE_DIR, ENV_NOTIFY_KAFKA_QUEUE_LIMIT, ]; diff --git a/crates/targets/src/check.rs b/crates/targets/src/check.rs index 2054f5cc4..c0693ea06 100644 --- a/crates/targets/src/check.rs +++ b/crates/targets/src/check.rs @@ -228,9 +228,11 @@ pub async fn check_postgres_server_available(args: &crate::target::postgres::Pos pub async fn check_kafka_broker_available(args: &crate::target::kafka::KafkaArgs) -> Result<(), crate::TargetError> { use rustfs_kafka_async::error::{ConnectionError, Error as KafkaError}; - use rustfs_kafka_async::{AsyncProducer, AsyncProducerConfig, RequiredAcks, SecurityConfig}; + use rustfs_kafka_async::{AsyncProducer, AsyncProducerConfig, RequiredAcks}; use std::time::Duration; + args.validate()?; + let map_kafka_error = |err: KafkaError, context: &str| match err { KafkaError::Connection(ConnectionError::NoHostReachable) => crate::TargetError::NotConnected, KafkaError::Connection(ConnectionError::Timeout(_)) => crate::TargetError::Timeout(format!("{context}: {err}")), @@ -249,14 +251,7 @@ pub async fn check_kafka_broker_available(args: &crate::target::kafka::KafkaArgs .with_ack_timeout(Duration::from_secs(5)) .with_required_acks(acks); - if args.tls_enable { - let mut security = SecurityConfig::new(); - if !args.tls_ca.is_empty() { - security = security.with_ca_cert(args.tls_ca.clone()); - } - if !args.tls_client_cert.is_empty() && !args.tls_client_key.is_empty() { - security = security.with_client_cert(args.tls_client_cert.clone(), args.tls_client_key.clone()); - } + if let Some(security) = args.security_config(false)? { config = config.with_security(security); } diff --git a/crates/targets/src/config/target_args.rs b/crates/targets/src/config/target_args.rs index 93b1825b2..57af9c875 100644 --- a/crates/targets/src/config/target_args.rs +++ b/crates/targets/src/config/target_args.rs @@ -17,7 +17,7 @@ use crate::error::TargetError; use crate::target::{ TargetType, amqp::AMQPArgs, - kafka::KafkaArgs, + kafka::{KAFKA_SASL_PLAIN, KafkaArgs}, mqtt::{MQTTArgs, MQTTTlsConfig, validate_mqtt_broker_url}, mysql::MySqlArgs, nats::{NATSArgs, validate_nats_address}, @@ -30,21 +30,22 @@ use rumqttc::QoS; use rustfs_config::{ AMQP_EXCHANGE, AMQP_MANDATORY, AMQP_PASSWORD, AMQP_PERSISTENT, AMQP_QUEUE_DIR, AMQP_QUEUE_LIMIT, AMQP_ROUTING_KEY, AMQP_TLS_CA, AMQP_TLS_CLIENT_CERT, AMQP_TLS_CLIENT_KEY, AMQP_URL, AMQP_USERNAME, DEFAULT_LIMIT, KAFKA_ACKS, KAFKA_BROKERS, - KAFKA_QUEUE_DIR, KAFKA_QUEUE_LIMIT, KAFKA_TLS_CA, KAFKA_TLS_CLIENT_CERT, KAFKA_TLS_CLIENT_KEY, KAFKA_TLS_ENABLE, KAFKA_TOPIC, - MQTT_BROKER, MQTT_KEEP_ALIVE_INTERVAL, MQTT_PASSWORD, MQTT_QOS, MQTT_QUEUE_DIR, MQTT_QUEUE_LIMIT, MQTT_RECONNECT_INTERVAL, - MQTT_TLS_CA, MQTT_TLS_CLIENT_CERT, MQTT_TLS_CLIENT_KEY, MQTT_TLS_POLICY, MQTT_TLS_TRUST_LEAF_AS_CA, MQTT_TOPIC, - MQTT_USERNAME, MQTT_WS_PATH_ALLOWLIST, MYSQL_DSN_STRING, MYSQL_FORMAT, MYSQL_MAX_OPEN_CONNECTIONS, MYSQL_QUEUE_DIR, - MYSQL_QUEUE_LIMIT, MYSQL_TABLE, MYSQL_TLS_CA, MYSQL_TLS_CLIENT_CERT, MYSQL_TLS_CLIENT_KEY, NATS_ADDRESS, - NATS_CREDENTIALS_FILE, NATS_PASSWORD, NATS_QUEUE_DIR, NATS_QUEUE_LIMIT, NATS_SUBJECT, NATS_TLS_CA, NATS_TLS_CLIENT_CERT, - NATS_TLS_CLIENT_KEY, NATS_TLS_REQUIRED, NATS_TOKEN, NATS_USERNAME, POSTGRES_DSN_STRING, POSTGRES_FORMAT, POSTGRES_QUEUE_DIR, - POSTGRES_QUEUE_LIMIT, POSTGRES_TABLE, POSTGRES_TLS_CA, POSTGRES_TLS_CLIENT_CERT, POSTGRES_TLS_CLIENT_KEY, - POSTGRES_TLS_REQUIRED, PULSAR_AUTH_TOKEN, PULSAR_BROKER, PULSAR_PASSWORD, PULSAR_QUEUE_DIR, PULSAR_QUEUE_LIMIT, - PULSAR_TLS_ALLOW_INSECURE, PULSAR_TLS_CA, PULSAR_TLS_HOSTNAME_VERIFICATION, PULSAR_TOPIC, PULSAR_USERNAME, REDIS_CHANNEL, - REDIS_CONNECTION_TIMEOUT, REDIS_KEEP_ALIVE_INTERVAL, REDIS_MAX_RETRY_ATTEMPTS, REDIS_MAX_RETRY_DELAY, REDIS_MIN_RETRY_DELAY, - REDIS_PASSWORD, REDIS_PIPELINE_BUFFER_SIZE, REDIS_QUEUE_DIR, REDIS_QUEUE_LIMIT, REDIS_RECONNECT_RETRY_ATTEMPTS, - REDIS_RESPONSE_TIMEOUT, REDIS_TLS_ALLOW_INSECURE, REDIS_TLS_CA, REDIS_TLS_CLIENT_CERT, REDIS_TLS_CLIENT_KEY, - REDIS_TLS_POLICY, REDIS_URL, REDIS_USERNAME, RUSTFS_WEBHOOK_SKIP_TLS_VERIFY_DEFAULT, WEBHOOK_AUTH_TOKEN, WEBHOOK_CLIENT_CA, - WEBHOOK_CLIENT_CERT, WEBHOOK_CLIENT_KEY, WEBHOOK_ENDPOINT, WEBHOOK_QUEUE_DIR, WEBHOOK_QUEUE_LIMIT, WEBHOOK_SKIP_TLS_VERIFY, + KAFKA_QUEUE_DIR, KAFKA_QUEUE_LIMIT, KAFKA_SASL_ENABLE, KAFKA_SASL_MECHANISM, KAFKA_SASL_PASSWORD, KAFKA_SASL_USERNAME, + KAFKA_TLS_CA, KAFKA_TLS_CLIENT_CERT, KAFKA_TLS_CLIENT_KEY, KAFKA_TLS_ENABLE, KAFKA_TOPIC, MQTT_BROKER, + MQTT_KEEP_ALIVE_INTERVAL, MQTT_PASSWORD, MQTT_QOS, MQTT_QUEUE_DIR, MQTT_QUEUE_LIMIT, MQTT_RECONNECT_INTERVAL, MQTT_TLS_CA, + MQTT_TLS_CLIENT_CERT, MQTT_TLS_CLIENT_KEY, MQTT_TLS_POLICY, MQTT_TLS_TRUST_LEAF_AS_CA, MQTT_TOPIC, MQTT_USERNAME, + MQTT_WS_PATH_ALLOWLIST, MYSQL_DSN_STRING, MYSQL_FORMAT, MYSQL_MAX_OPEN_CONNECTIONS, MYSQL_QUEUE_DIR, MYSQL_QUEUE_LIMIT, + MYSQL_TABLE, MYSQL_TLS_CA, MYSQL_TLS_CLIENT_CERT, MYSQL_TLS_CLIENT_KEY, NATS_ADDRESS, NATS_CREDENTIALS_FILE, NATS_PASSWORD, + NATS_QUEUE_DIR, NATS_QUEUE_LIMIT, NATS_SUBJECT, NATS_TLS_CA, NATS_TLS_CLIENT_CERT, NATS_TLS_CLIENT_KEY, NATS_TLS_REQUIRED, + NATS_TOKEN, NATS_USERNAME, POSTGRES_DSN_STRING, POSTGRES_FORMAT, POSTGRES_QUEUE_DIR, POSTGRES_QUEUE_LIMIT, POSTGRES_TABLE, + POSTGRES_TLS_CA, POSTGRES_TLS_CLIENT_CERT, POSTGRES_TLS_CLIENT_KEY, POSTGRES_TLS_REQUIRED, PULSAR_AUTH_TOKEN, PULSAR_BROKER, + PULSAR_PASSWORD, PULSAR_QUEUE_DIR, PULSAR_QUEUE_LIMIT, PULSAR_TLS_ALLOW_INSECURE, PULSAR_TLS_CA, + PULSAR_TLS_HOSTNAME_VERIFICATION, PULSAR_TOPIC, PULSAR_USERNAME, REDIS_CHANNEL, REDIS_CONNECTION_TIMEOUT, + REDIS_KEEP_ALIVE_INTERVAL, REDIS_MAX_RETRY_ATTEMPTS, REDIS_MAX_RETRY_DELAY, REDIS_MIN_RETRY_DELAY, REDIS_PASSWORD, + REDIS_PIPELINE_BUFFER_SIZE, REDIS_QUEUE_DIR, REDIS_QUEUE_LIMIT, REDIS_RECONNECT_RETRY_ATTEMPTS, REDIS_RESPONSE_TIMEOUT, + REDIS_TLS_ALLOW_INSECURE, REDIS_TLS_CA, REDIS_TLS_CLIENT_CERT, REDIS_TLS_CLIENT_KEY, REDIS_TLS_POLICY, REDIS_URL, + REDIS_USERNAME, RUSTFS_WEBHOOK_SKIP_TLS_VERIFY_DEFAULT, WEBHOOK_AUTH_TOKEN, WEBHOOK_CLIENT_CA, WEBHOOK_CLIENT_CERT, + WEBHOOK_CLIENT_KEY, WEBHOOK_ENDPOINT, WEBHOOK_QUEUE_DIR, WEBHOOK_QUEUE_LIMIT, WEBHOOK_SKIP_TLS_VERIFY, }; use rustfs_ecstore::config::KVS; use std::path::Path; @@ -68,6 +69,19 @@ fn parse_kafka_acks_value(value: Option<&str>) -> Result { } } +fn parse_kafka_sasl_enable(config: &KVS, has_sasl_fields: bool) -> Result { + match config.lookup(KAFKA_SASL_ENABLE) { + Some(value) => { + if value.trim().is_empty() { + return Ok(has_sasl_fields); + } + parse_target_bool(Some(value.as_str())) + .ok_or_else(|| TargetError::Configuration(format!("Invalid Kafka {KAFKA_SASL_ENABLE} boolean value: {value}"))) + } + None => Ok(has_sasl_fields), + } +} + fn parse_amqp_bool_value(field: &str, config: &KVS, default: bool) -> Result { match config.lookup(field) { Some(value) => parse_target_bool(Some(value.as_str())) @@ -482,7 +496,13 @@ pub fn build_kafka_args(config: &KVS, default_queue_dir: &str, target_type: Targ .lookup(KAFKA_TOPIC) .ok_or_else(|| TargetError::Configuration("Missing Kafka topic".to_string()))?; - Ok(KafkaArgs { + let sasl_mechanism = config.lookup(KAFKA_SASL_MECHANISM).unwrap_or_default(); + let sasl_username = config.lookup(KAFKA_SASL_USERNAME).unwrap_or_default(); + let sasl_password = config.lookup(KAFKA_SASL_PASSWORD).unwrap_or_default(); + let has_sasl_fields = !sasl_mechanism.trim().is_empty() || !sasl_username.is_empty() || !sasl_password.is_empty(); + let sasl_enable = parse_kafka_sasl_enable(config, has_sasl_fields)?; + + let args = KafkaArgs { enable: true, brokers, topic, @@ -491,6 +511,14 @@ pub fn build_kafka_args(config: &KVS, default_queue_dir: &str, target_type: Targ tls_ca: config.lookup(KAFKA_TLS_CA).unwrap_or_default(), tls_client_cert: config.lookup(KAFKA_TLS_CLIENT_CERT).unwrap_or_default(), tls_client_key: config.lookup(KAFKA_TLS_CLIENT_KEY).unwrap_or_default(), + sasl_enable, + sasl_mechanism: if sasl_enable && sasl_mechanism.trim().is_empty() { + KAFKA_SASL_PLAIN.to_string() + } else { + sasl_mechanism.trim().to_string() + }, + sasl_username, + sasl_password, queue_dir: config .lookup(KAFKA_QUEUE_DIR) .unwrap_or_else(|| default_queue_dir.to_string()), @@ -499,38 +527,13 @@ pub fn build_kafka_args(config: &KVS, default_queue_dir: &str, target_type: Targ .and_then(|v| v.parse::().ok()) .unwrap_or(DEFAULT_LIMIT), target_type, - }) + }; + args.validate()?; + Ok(args) } pub fn validate_kafka_config(config: &KVS, default_queue_dir: &str) -> Result<(), TargetError> { - let brokers_raw = config - .lookup(KAFKA_BROKERS) - .ok_or_else(|| TargetError::Configuration("Missing Kafka brokers".to_string()))?; - if brokers_raw.split(',').map(|s| s.trim()).all(|s| s.is_empty()) { - return Err(TargetError::Configuration("Kafka brokers cannot be empty".to_string())); - } - - if config.lookup(KAFKA_TOPIC).is_none() { - return Err(TargetError::Configuration("Missing Kafka topic".to_string())); - } - - parse_kafka_acks_value(config.lookup(KAFKA_ACKS).as_deref())?; - - let tls_client_cert = config.lookup(KAFKA_TLS_CLIENT_CERT).unwrap_or_default(); - let tls_client_key = config.lookup(KAFKA_TLS_CLIENT_KEY).unwrap_or_default(); - if tls_client_cert.is_empty() != tls_client_key.is_empty() { - return Err(TargetError::Configuration( - "Kafka tls_client_cert and tls_client_key must be specified together".to_string(), - )); - } - - let queue_dir = config - .lookup(KAFKA_QUEUE_DIR) - .unwrap_or_else(|| default_queue_dir.to_string()); - if !queue_dir.is_empty() && !Path::new(&queue_dir).is_absolute() { - return Err(TargetError::Configuration("Kafka queue directory must be an absolute path".to_string())); - } - + let _ = build_kafka_args(config, default_queue_dir, TargetType::NotifyEvent)?; Ok(()) } @@ -594,14 +597,19 @@ mod tests { build_amqp_args, build_kafka_args, build_mysql_args, build_postgres_args, build_redis_args, validate_amqp_config, validate_kafka_config, validate_mysql_config, validate_postgres_config, validate_redis_config, }; - use crate::target::{TargetType, postgres::PostgresFormat}; + use crate::target::{ + TargetType, + kafka::{KAFKA_SASL_PLAIN, KAFKA_SASL_SCRAM_SHA_512}, + postgres::PostgresFormat, + }; use rustfs_config::{ AMQP_EXCHANGE, AMQP_MANDATORY, AMQP_PASSWORD, AMQP_PERSISTENT, AMQP_QUEUE_DIR, AMQP_ROUTING_KEY, AMQP_TLS_CLIENT_CERT, - AMQP_TLS_CLIENT_KEY, AMQP_URL, AMQP_USERNAME, KAFKA_ACKS, KAFKA_BROKERS, KAFKA_TOPIC, MYSQL_DSN_STRING, - MYSQL_MAX_OPEN_CONNECTIONS, MYSQL_QUEUE_DIR, MYSQL_TABLE, MYSQL_TLS_CA, MYSQL_TLS_CLIENT_CERT, MYSQL_TLS_CLIENT_KEY, - POSTGRES_DSN_STRING, POSTGRES_FORMAT, POSTGRES_QUEUE_DIR, POSTGRES_TABLE, POSTGRES_TLS_CA, POSTGRES_TLS_CLIENT_CERT, - POSTGRES_TLS_CLIENT_KEY, REDIS_CHANNEL, REDIS_CONNECTION_TIMEOUT, REDIS_MAX_RETRY_DELAY, REDIS_MIN_RETRY_DELAY, - REDIS_PIPELINE_BUFFER_SIZE, REDIS_RECONNECT_RETRY_ATTEMPTS, REDIS_RESPONSE_TIMEOUT, REDIS_TLS_ALLOW_INSECURE, REDIS_URL, + AMQP_TLS_CLIENT_KEY, AMQP_URL, AMQP_USERNAME, KAFKA_ACKS, KAFKA_BROKERS, KAFKA_SASL_ENABLE, KAFKA_SASL_MECHANISM, + KAFKA_SASL_PASSWORD, KAFKA_SASL_USERNAME, KAFKA_TLS_ENABLE, KAFKA_TOPIC, MYSQL_DSN_STRING, MYSQL_MAX_OPEN_CONNECTIONS, + MYSQL_QUEUE_DIR, MYSQL_TABLE, MYSQL_TLS_CA, MYSQL_TLS_CLIENT_CERT, MYSQL_TLS_CLIENT_KEY, POSTGRES_DSN_STRING, + POSTGRES_FORMAT, POSTGRES_QUEUE_DIR, POSTGRES_TABLE, POSTGRES_TLS_CA, POSTGRES_TLS_CLIENT_CERT, POSTGRES_TLS_CLIENT_KEY, + REDIS_CHANNEL, REDIS_CONNECTION_TIMEOUT, REDIS_MAX_RETRY_DELAY, REDIS_MIN_RETRY_DELAY, REDIS_PIPELINE_BUFFER_SIZE, + REDIS_RECONNECT_RETRY_ATTEMPTS, REDIS_RESPONSE_TIMEOUT, REDIS_TLS_ALLOW_INSECURE, REDIS_URL, }; use rustfs_ecstore::config::KVS; @@ -778,6 +786,75 @@ mod tests { assert!(err.to_string().contains("Kafka acks must be one of")); } + #[test] + fn build_kafka_args_infers_sasl_enable_from_credentials() { + let mut config = kafka_base_config(); + config.insert(KAFKA_TLS_ENABLE.to_string(), "on".to_string()); + config.insert(KAFKA_SASL_USERNAME.to_string(), "user".to_string()); + config.insert(KAFKA_SASL_PASSWORD.to_string(), "secret".to_string()); + + let args = build_kafka_args(&config, "", TargetType::NotifyEvent).expect("valid kafka SASL args"); + + assert!(args.sasl_enable); + assert_eq!(args.sasl_mechanism, KAFKA_SASL_PLAIN); + assert_eq!(args.sasl_username, "user"); + assert_eq!(args.sasl_password, "secret"); + } + + #[test] + fn build_kafka_args_infers_sasl_enable_when_enable_is_blank() { + let mut config = kafka_base_config(); + config.insert(KAFKA_TLS_ENABLE.to_string(), "on".to_string()); + config.insert(KAFKA_SASL_ENABLE.to_string(), "".to_string()); + config.insert(KAFKA_SASL_USERNAME.to_string(), "user".to_string()); + config.insert(KAFKA_SASL_PASSWORD.to_string(), "secret".to_string()); + + let args = + build_kafka_args(&config, "", TargetType::NotifyEvent).expect("blank sasl_enable should infer from SASL fields"); + + assert!(args.sasl_enable); + assert_eq!(args.sasl_mechanism, KAFKA_SASL_PLAIN); + } + + #[test] + fn build_kafka_args_accepts_scram_sha_512() { + let mut config = kafka_base_config(); + config.insert(KAFKA_TLS_ENABLE.to_string(), "true".to_string()); + config.insert(KAFKA_SASL_ENABLE.to_string(), "true".to_string()); + config.insert(KAFKA_SASL_MECHANISM.to_string(), "SCRAM-SHA-512".to_string()); + config.insert(KAFKA_SASL_USERNAME.to_string(), "user".to_string()); + config.insert(KAFKA_SASL_PASSWORD.to_string(), "secret".to_string()); + + let args = build_kafka_args(&config, "", TargetType::NotifyEvent).expect("valid kafka SCRAM args"); + + assert!(args.sasl_enable); + assert_eq!(args.sasl_mechanism, KAFKA_SASL_SCRAM_SHA_512); + } + + #[test] + fn validate_kafka_config_rejects_sasl_without_tls() { + let mut config = kafka_base_config(); + config.insert(KAFKA_SASL_ENABLE.to_string(), "on".to_string()); + config.insert(KAFKA_SASL_USERNAME.to_string(), "user".to_string()); + config.insert(KAFKA_SASL_PASSWORD.to_string(), "secret".to_string()); + + let err = validate_kafka_config(&config, "").expect_err("SASL without TLS should fail"); + + assert!(err.to_string().contains("requires tls_enable")); + } + + #[test] + fn validate_kafka_config_rejects_sasl_fields_when_disabled() { + let mut config = kafka_base_config(); + config.insert(KAFKA_SASL_ENABLE.to_string(), "off".to_string()); + config.insert(KAFKA_SASL_USERNAME.to_string(), "user".to_string()); + config.insert(KAFKA_SASL_PASSWORD.to_string(), "secret".to_string()); + + let err = validate_kafka_config(&config, "").expect_err("SASL fields with disabled SASL should fail"); + + assert!(err.to_string().contains("sasl_enable must be true")); + } + #[test] fn build_mysql_args_accepts_minimal_config() { let args = build_mysql_args(&mysql_base_config(), "", TargetType::NotifyEvent).expect("valid mysql args"); diff --git a/crates/targets/src/manifest.rs b/crates/targets/src/manifest.rs index d858b1092..6a683197f 100644 --- a/crates/targets/src/manifest.rs +++ b/crates/targets/src/manifest.rs @@ -14,8 +14,8 @@ use crate::domain::TargetDomain; use rustfs_config::{ - AMQP_PASSWORD, AMQP_TLS_CLIENT_CERT, AMQP_TLS_CLIENT_KEY, KAFKA_TLS_CLIENT_CERT, KAFKA_TLS_CLIENT_KEY, MQTT_PASSWORD, - MQTT_TLS_CLIENT_CERT, MQTT_TLS_CLIENT_KEY, MYSQL_DSN_STRING, MYSQL_TLS_CLIENT_CERT, MYSQL_TLS_CLIENT_KEY, + AMQP_PASSWORD, AMQP_TLS_CLIENT_CERT, AMQP_TLS_CLIENT_KEY, KAFKA_SASL_PASSWORD, KAFKA_TLS_CLIENT_CERT, KAFKA_TLS_CLIENT_KEY, + MQTT_PASSWORD, MQTT_TLS_CLIENT_CERT, MQTT_TLS_CLIENT_KEY, MYSQL_DSN_STRING, MYSQL_TLS_CLIENT_CERT, MYSQL_TLS_CLIENT_KEY, NATS_CREDENTIALS_FILE, NATS_PASSWORD, NATS_TLS_CLIENT_CERT, NATS_TLS_CLIENT_KEY, NATS_TOKEN, POSTGRES_DSN_STRING, POSTGRES_TLS_CLIENT_CERT, POSTGRES_TLS_CLIENT_KEY, PULSAR_AUTH_TOKEN, PULSAR_PASSWORD, REDIS_PASSWORD, REDIS_TLS_CLIENT_CERT, REDIS_TLS_CLIENT_KEY, WEBHOOK_AUTH_TOKEN, WEBHOOK_CLIENT_CERT, WEBHOOK_CLIENT_KEY, @@ -106,7 +106,7 @@ const NO_SECRET_FIELDS: &[&str] = &[]; const WEBHOOK_SECRET_FIELDS: &[&str] = &[WEBHOOK_AUTH_TOKEN, WEBHOOK_CLIENT_CERT, WEBHOOK_CLIENT_KEY]; const MQTT_SECRET_FIELDS: &[&str] = &[MQTT_PASSWORD, MQTT_TLS_CLIENT_CERT, MQTT_TLS_CLIENT_KEY]; -const KAFKA_SECRET_FIELDS: &[&str] = &[KAFKA_TLS_CLIENT_CERT, KAFKA_TLS_CLIENT_KEY]; +const KAFKA_SECRET_FIELDS: &[&str] = &[KAFKA_TLS_CLIENT_CERT, KAFKA_TLS_CLIENT_KEY, KAFKA_SASL_PASSWORD]; const AMQP_SECRET_FIELDS: &[&str] = &[AMQP_PASSWORD, AMQP_TLS_CLIENT_CERT, AMQP_TLS_CLIENT_KEY]; const NATS_SECRET_FIELDS: &[&str] = &[ NATS_PASSWORD, @@ -221,7 +221,7 @@ mod tests { installable_target_marketplace_manifest, }; use crate::domain::TargetDomain; - use rustfs_config::{WEBHOOK_AUTH_TOKEN, WEBHOOK_CLIENT_CERT, WEBHOOK_CLIENT_KEY}; + use rustfs_config::{KAFKA_SASL_PASSWORD, WEBHOOK_AUTH_TOKEN, WEBHOOK_CLIENT_CERT, WEBHOOK_CLIENT_KEY}; #[test] fn builtin_webhook_manifest_marks_secret_fields() { @@ -261,6 +261,13 @@ mod tests { assert_eq!(manifest.supported_domains, &[TargetDomain::Audit, TargetDomain::Notify]); } + #[test] + fn builtin_kafka_manifest_marks_sasl_password_secret() { + let manifest = builtin_target_manifest("kafka"); + + assert!(manifest.secret_fields.contains(&KAFKA_SASL_PASSWORD)); + } + #[test] fn marketplace_manifest_from_builtin_manifest_is_stable() { let base = builtin_target_manifest("redis"); diff --git a/crates/targets/src/target/kafka.rs b/crates/targets/src/target/kafka.rs index dcd3521c4..8baff55a2 100644 --- a/crates/targets/src/target/kafka.rs +++ b/crates/targets/src/target/kafka.rs @@ -29,16 +29,20 @@ use crate::{ }; use async_trait::async_trait; use rustfs_kafka_async::error::{ConnectionError, Error as KafkaError}; -use rustfs_kafka_async::{AsyncProducer, AsyncProducerConfig, Record, RequiredAcks, SecurityConfig}; +use rustfs_kafka_async::{AsyncProducer, AsyncProducerConfig, Record, RequiredAcks, SaslConfig, SecurityConfig}; use rustfs_tls_runtime::{load_cert_bundle_der_bytes, load_private_key}; use serde::Serialize; use serde::de::DeserializeOwned; -use std::{marker::PhantomData, sync::Arc, time::Duration}; +use std::{fmt, marker::PhantomData, sync::Arc, time::Duration}; use tokio::sync::Mutex; use tracing::{debug, error, info, instrument, warn}; +pub(crate) const KAFKA_SASL_PLAIN: &str = "PLAIN"; +pub(crate) const KAFKA_SASL_SCRAM_SHA_256: &str = "SCRAM-SHA-256"; +pub(crate) const KAFKA_SASL_SCRAM_SHA_512: &str = "SCRAM-SHA-512"; + /// Arguments for configuring a Kafka target -#[derive(Debug, Clone)] +#[derive(Clone)] pub struct KafkaArgs { /// Whether the target is enabled pub enable: bool, @@ -56,6 +60,14 @@ pub struct KafkaArgs { pub tls_client_cert: String, /// Optional path to client private key for mTLS pub tls_client_key: String, + /// Whether to enable SASL authentication over the TLS transport + pub sasl_enable: bool, + /// SASL mechanism (PLAIN, SCRAM-SHA-256, or SCRAM-SHA-512) + pub sasl_mechanism: String, + /// SASL username + pub sasl_username: String, + /// SASL password + pub sasl_password: String, /// The directory to store events in case of failure pub queue_dir: String, /// The maximum number of events to store @@ -64,6 +76,58 @@ pub struct KafkaArgs { pub target_type: TargetType, } +impl fmt::Debug for KafkaArgs { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("KafkaArgs") + .field("enable", &self.enable) + .field("brokers", &self.brokers) + .field("topic", &self.topic) + .field("acks", &self.acks) + .field("tls_enable", &self.tls_enable) + .field("tls_ca", &self.tls_ca) + .field("tls_client_cert", &self.tls_client_cert) + .field( + "tls_client_key", + if self.tls_client_key.is_empty() { + &"" + } else { + &"***REDACTED***" + }, + ) + .field("sasl_enable", &self.sasl_enable) + .field("sasl_mechanism", &self.sasl_mechanism) + .field("sasl_username", &self.sasl_username) + .field( + "sasl_password", + if self.sasl_password.is_empty() { + &"" + } else { + &"***REDACTED***" + }, + ) + .field("queue_dir", &self.queue_dir) + .field("queue_limit", &self.queue_limit) + .field("target_type", &self.target_type) + .finish() + } +} + +fn normalize_kafka_sasl_mechanism(mechanism: &str) -> Result<&'static str, TargetError> { + let mechanism = mechanism.trim(); + if mechanism.is_empty() || mechanism.eq_ignore_ascii_case(KAFKA_SASL_PLAIN) { + return Ok(KAFKA_SASL_PLAIN); + } + if mechanism.eq_ignore_ascii_case(KAFKA_SASL_SCRAM_SHA_256) { + return Ok(KAFKA_SASL_SCRAM_SHA_256); + } + if mechanism.eq_ignore_ascii_case(KAFKA_SASL_SCRAM_SHA_512) { + return Ok(KAFKA_SASL_SCRAM_SHA_512); + } + Err(TargetError::Configuration( + "kafka sasl_mechanism must be one of: PLAIN, SCRAM-SHA-256, SCRAM-SHA-512".to_string(), + )) +} + impl KafkaArgs { /// Validates the KafkaArgs configuration pub fn validate(&self) -> Result<(), TargetError> { @@ -89,6 +153,24 @@ impl KafkaArgs { )); } + if self.sasl_enable { + if !self.tls_enable { + return Err(TargetError::Configuration( + "kafka sasl_enable requires tls_enable for SASL_SSL".to_string(), + )); + } + normalize_kafka_sasl_mechanism(&self.sasl_mechanism)?; + if self.sasl_username.is_empty() || self.sasl_password.is_empty() { + return Err(TargetError::Configuration( + "kafka sasl_username and sasl_password must be specified when sasl_enable is true".to_string(), + )); + } + } else if !self.sasl_mechanism.is_empty() || !self.sasl_username.is_empty() || !self.sasl_password.is_empty() { + return Err(TargetError::Configuration( + "kafka sasl_enable must be true when SASL fields are specified".to_string(), + )); + } + if !self.queue_dir.is_empty() { let path = std::path::Path::new(&self.queue_dir); if !path.is_absolute() { @@ -98,6 +180,49 @@ impl KafkaArgs { Ok(()) } + + pub(crate) fn security_config(&self, validate_tls_files: bool) -> Result, TargetError> { + if !self.tls_enable && !self.sasl_enable { + return Ok(None); + } + + let mut security = SecurityConfig::new(); + if !self.tls_ca.is_empty() { + if validate_tls_files { + let certs = load_cert_bundle_der_bytes(&self.tls_ca) + .map_err(|e| TargetError::Configuration(format!("Failed to parse Kafka tls_ca: {e}")))?; + if certs.is_empty() { + return Err(TargetError::Configuration( + "Kafka tls_ca did not contain any parsable certificates".to_string(), + )); + } + } + security = security.with_ca_cert(self.tls_ca.clone()); + } + if !self.tls_client_cert.is_empty() && !self.tls_client_key.is_empty() { + if validate_tls_files { + let certs = load_cert_bundle_der_bytes(&self.tls_client_cert) + .map_err(|e| TargetError::Configuration(format!("Failed to parse Kafka tls_client_cert: {e}")))?; + if certs.is_empty() { + return Err(TargetError::Configuration( + "Kafka tls_client_cert did not contain any parsable certificates".to_string(), + )); + } + let _ = load_private_key(&self.tls_client_key) + .map_err(|e| TargetError::Configuration(format!("Failed to parse Kafka tls_client_key: {e}")))?; + } + security = security.with_client_cert(self.tls_client_cert.clone(), self.tls_client_key.clone()); + } + if self.sasl_enable { + security = security.with_sasl(SaslConfig::new( + normalize_kafka_sasl_mechanism(&self.sasl_mechanism)?.to_string(), + self.sasl_username.clone(), + self.sasl_password.clone(), + )); + } + + Ok(Some(security)) + } } /// A target that sends events to an Apache Kafka topic @@ -173,32 +298,7 @@ where .with_ack_timeout(Duration::from_secs(30)) .with_required_acks(acks); - if self.args.tls_enable { - let mut security = SecurityConfig::new(); - if !self.args.tls_ca.is_empty() { - let certs = load_cert_bundle_der_bytes(&self.args.tls_ca) - .map_err(|e| Self::map_kafka_error(KafkaError::Config(e.to_string()), "Failed to parse Kafka tls_ca"))?; - if certs.is_empty() { - return Err(TargetError::Configuration( - "Kafka tls_ca did not contain any parsable certificates".to_string(), - )); - } - security = security.with_ca_cert(self.args.tls_ca.clone()); - } - if !self.args.tls_client_cert.is_empty() && !self.args.tls_client_key.is_empty() { - let certs = load_cert_bundle_der_bytes(&self.args.tls_client_cert).map_err(|e| { - Self::map_kafka_error(KafkaError::Config(e.to_string()), "Failed to parse Kafka tls_client_cert") - })?; - if certs.is_empty() { - return Err(TargetError::Configuration( - "Kafka tls_client_cert did not contain any parsable certificates".to_string(), - )); - } - let _ = load_private_key(&self.args.tls_client_key).map_err(|e| { - Self::map_kafka_error(KafkaError::Config(e.to_string()), "Failed to parse Kafka tls_client_key") - })?; - security = security.with_client_cert(self.args.tls_client_cert.clone(), self.args.tls_client_key.clone()); - } + if let Some(security) = self.args.security_config(true)? { config = config.with_security(security); } @@ -438,6 +538,10 @@ mod tests { tls_ca: String::new(), tls_client_cert: String::new(), tls_client_key: String::new(), + sasl_enable: false, + sasl_mechanism: String::new(), + sasl_username: String::new(), + sasl_password: String::new(), queue_dir: String::new(), queue_limit: 0, target_type: TargetType::NotifyEvent, @@ -496,4 +600,84 @@ mod tests { }; assert!(args.validate().is_err()); } + + #[test] + fn test_validate_sasl_requires_tls() { + let args = KafkaArgs { + sasl_enable: true, + sasl_mechanism: KAFKA_SASL_SCRAM_SHA_512.to_string(), + sasl_username: "user".to_string(), + sasl_password: "secret".to_string(), + ..base_args() + }; + let err = args.validate().expect_err("SASL without TLS should fail"); + assert!(err.to_string().contains("requires tls_enable")); + } + + #[test] + fn test_validate_sasl_requires_username_and_password() { + let args = KafkaArgs { + tls_enable: true, + sasl_enable: true, + sasl_mechanism: KAFKA_SASL_PLAIN.to_string(), + sasl_username: "user".to_string(), + sasl_password: String::new(), + ..base_args() + }; + let err = args.validate().expect_err("SASL credentials should be paired"); + assert!(err.to_string().contains("sasl_username and sasl_password")); + } + + #[test] + fn test_validate_sasl_rejects_unsupported_mechanism() { + let args = KafkaArgs { + tls_enable: true, + sasl_enable: true, + sasl_mechanism: "OAUTHBEARER".to_string(), + sasl_username: "user".to_string(), + sasl_password: "secret".to_string(), + ..base_args() + }; + let err = args.validate().expect_err("unsupported SASL mechanism should fail"); + assert!(err.to_string().contains("sasl_mechanism must be one of")); + } + + #[test] + fn test_security_config_includes_sasl() { + let args = KafkaArgs { + tls_enable: true, + sasl_enable: true, + sasl_mechanism: "scram-sha-512".to_string(), + sasl_username: "user".to_string(), + sasl_password: "secret".to_string(), + ..base_args() + }; + + let security = args + .security_config(false) + .expect("valid security config") + .expect("security should be configured"); + let sasl = security.sasl_config().expect("SASL should be configured"); + + assert_eq!(sasl.mechanism(), KAFKA_SASL_SCRAM_SHA_512); + assert_eq!(sasl.username(), "user"); + assert_eq!(sasl.password(), "secret"); + } + + #[test] + fn test_debug_redacts_sasl_password_and_tls_key() { + let rendered = format!( + "{:?}", + KafkaArgs { + tls_client_key: "/tmp/client.key".to_string(), + sasl_enable: true, + sasl_password: "super-secret".to_string(), + ..base_args() + } + ); + + assert!(!rendered.contains("super-secret")); + assert!(!rendered.contains("/tmp/client.key")); + assert!(rendered.contains("***REDACTED***")); + } }