From 9ce3c7742cf37ce888c56acffae1829fecf28c24 Mon Sep 17 00:00:00 2001 From: houseme Date: Tue, 17 Mar 2026 23:53:28 +0800 Subject: [PATCH] fix(targets): pass credentials to MQTT broker availability check (#2192) Signed-off-by: houseme Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- crates/targets/src/check.rs | 25 +++++++++++++++++++++++-- rustfs/src/admin/handlers/event.rs | 4 +++- 2 files changed, 26 insertions(+), 3 deletions(-) diff --git a/crates/targets/src/check.rs b/crates/targets/src/check.rs index 5e6ab53d3..fcc85b507 100644 --- a/crates/targets/src/check.rs +++ b/crates/targets/src/check.rs @@ -17,6 +17,8 @@ /// # Arguments /// * `broker_url` - URL of MQTT Broker, for example `mqtt://localhost:1883` /// * `topic` - Topic for testing connections +/// * `username` - Optional username for authentication +/// * `password` - Optional password for authentication /// # Returns /// * `Ok(())` - If the connection is successful /// * `Err(String)` - If the connection fails, contains an error message @@ -25,7 +27,12 @@ /// ```rust,no_run /// #[tokio::main] /// async fn main() { -/// let result = rustfs_targets::check_mqtt_broker_available("mqtt://localhost:1883", "test/topic").await; +/// let result = rustfs_targets::check_mqtt_broker_available( +/// "mqtt://localhost:1883", +/// "test/topic", +/// Some("myuser"), +/// Some("mypass"), +/// ).await; /// if result.is_ok() { /// println!("MQTT Broker is available"); /// } else { @@ -42,7 +49,12 @@ /// tokio = { version = "1", features = ["full"] } /// ``` /// -pub async fn check_mqtt_broker_available(broker_url: &str, topic: &str) -> Result<(), String> { +pub async fn check_mqtt_broker_available( + broker_url: &str, + topic: &str, + username: Option<&str>, + password: Option<&str>, +) -> Result<(), String> { use rumqttc::{AsyncClient, MqttOptions, QoS}; let url = rustfs_utils::parse_url(broker_url).map_err(|e| format!("Broker URL parsing failed:{e}"))?; let url = url.url(); @@ -55,6 +67,15 @@ pub async fn check_mqtt_broker_available(broker_url: &str, topic: &str) -> Resul let host = url.host_str().ok_or("Broker is missing host")?; let port = url.port().unwrap_or(1883); let mut mqtt_options = MqttOptions::new("rustfs_check", host, port); + + // Set credentials if provided + if let Some(user) = username + && !user.is_empty() + { + let pass = password.unwrap_or(""); + mqtt_options.set_credentials(user, pass); + } + mqtt_options.set_keep_alive(std::time::Duration::from_secs(5)); let (client, mut eventloop) = AsyncClient::new(mqtt_options, 1); diff --git a/rustfs/src/admin/handlers/event.rs b/rustfs/src/admin/handlers/event.rs index 01d9e9f60..4d0c879ee 100644 --- a/rustfs/src/admin/handlers/event.rs +++ b/rustfs/src/admin/handlers/event.rs @@ -225,7 +225,9 @@ impl Operation for NotificationTarget { let topic = kv_map .get(rustfs_config::MQTT_TOPIC) .ok_or_else(|| s3_error!(InvalidArgument, "topic is required"))?; - check_mqtt_broker_available(endpoint, topic) + let username = kv_map.get(rustfs_config::MQTT_USERNAME).copied(); + let password = kv_map.get(rustfs_config::MQTT_PASSWORD).copied(); + check_mqtt_broker_available(endpoint, topic, username, password) .await .map_err(|e| s3_error!(InvalidArgument, "MQTT Broker unavailable: {}", e))?;