diff --git a/crates/audit/src/factory.rs b/crates/audit/src/factory.rs index 9ee36731e..0b51fcc07 100644 --- a/crates/audit/src/factory.rs +++ b/crates/audit/src/factory.rs @@ -120,7 +120,7 @@ pub fn builtin_target_descriptors() -> Vec> ), BuiltinTargetDescriptor::new( rustfs_config::audit::AUDIT_MYSQL_SUB_SYS, - TargetRequestValidator::MySql, + TargetRequestValidator::MySql(TargetType::AuditLog), TargetPluginDescriptor::new( ChannelTargetType::MySql.as_str(), AUDIT_MYSQL_KEYS, diff --git a/crates/notify/src/factory.rs b/crates/notify/src/factory.rs index 6489f6053..a948cb344 100644 --- a/crates/notify/src/factory.rs +++ b/crates/notify/src/factory.rs @@ -85,7 +85,7 @@ pub fn builtin_target_descriptors() -> Vec> { ), BuiltinTargetDescriptor::new( NOTIFY_MYSQL_SUB_SYS, - TargetRequestValidator::MySql, + TargetRequestValidator::MySql(TargetType::NotifyEvent), TargetPluginDescriptor::new( ChannelTargetType::MySql.as_str(), NOTIFY_MYSQL_KEYS, diff --git a/crates/targets/src/check.rs b/crates/targets/src/check.rs index 7b486137f..1d0e949fa 100644 --- a/crates/targets/src/check.rs +++ b/crates/targets/src/check.rs @@ -123,6 +123,64 @@ pub async fn check_pulsar_broker_available(args: &crate::target::pulsar::PulsarA .unwrap_or_else(|_| Err(crate::TargetError::Timeout("Pulsar connection timed out".to_string()))) } +/// Probes a MySQL server for connectivity. +/// +/// 1. Validates `args`. +/// 2. Parses the DSN and builds a connection pool. +/// 3. Runs `SELECT 1` to confirm credentials work. +pub async fn check_mysql_server_available(args: &crate::target::mysql::MySqlArgs) -> Result<(), crate::TargetError> { + use crate::target::ensure_rustls_provider_installed; + use crate::target::mysql::{MySqlDsn, map_mysql_error}; + use mysql_async::{Opts, OptsBuilder, Pool, SslOpts, prelude::Queryable}; + use std::path::PathBuf; + + args.validate()?; + + let dsn = MySqlDsn::parse(&args.dsn_string)?; + + let mut builder = OptsBuilder::default() + .user(Some(dsn.user.clone())) + .pass(Some(dsn.password.clone())) + .ip_or_hostname(dsn.host.clone()) + .tcp_port(dsn.port) + .db_name(Some(dsn.database.clone())); + + if dsn.tls { + ensure_rustls_provider_installed(); + let mut ssl_opts = SslOpts::default(); + if !args.tls_ca.is_empty() { + ssl_opts = ssl_opts.with_root_certs(vec![PathBuf::from(args.tls_ca.clone()).into()]); + } + if !args.tls_client_cert.is_empty() && !args.tls_client_key.is_empty() { + let identity = mysql_async::ClientIdentity::new( + PathBuf::from(args.tls_client_cert.clone()).into(), + PathBuf::from(args.tls_client_key.clone()).into(), + ); + ssl_opts = ssl_opts.with_client_identity(Some(identity)); + } + builder = builder.ssl_opts(Some(ssl_opts)); + } + + let pool = Pool::new(Opts::from(builder)); + // Pool is dropped at scope exit; pool.disconnect() is deliberately + // avoided — integration tests show it hangs indefinitely, exceeding + // the 8s timeout. Drops handle cleanup without blocking. + + let timeout = std::time::Duration::from_secs(8); + tokio::time::timeout(timeout, async { + let mut conn = pool + .get_conn() + .await + .map_err(|err| map_mysql_error(err, "MySQL connectivity probe failed to acquire connection"))?; + conn.query_drop("SELECT 1") + .await + .map_err(|err| map_mysql_error(err, "MySQL connectivity probe failed"))?; + Ok::<(), crate::TargetError>(()) + }) + .await + .unwrap_or_else(|_| Err(crate::TargetError::Timeout("MySQL connectivity probe timed out".to_string()))) +} + /// Probes a PostgreSQL server for connectivity and verifies the configured /// table is readable. /// diff --git a/crates/targets/src/lib.rs b/crates/targets/src/lib.rs index aeff41e37..a2c407e98 100644 --- a/crates/targets/src/lib.rs +++ b/crates/targets/src/lib.rs @@ -23,7 +23,8 @@ pub mod target; pub use check::{ check_amqp_broker_available, check_kafka_broker_available, check_mqtt_broker_available, check_mqtt_broker_available_with_tls, - check_nats_server_available, check_postgres_server_available, check_pulsar_broker_available, check_redis_server_available, + check_mysql_server_available, check_nats_server_available, check_postgres_server_available, check_pulsar_broker_available, + check_redis_server_available, }; pub use error::{StoreError, TargetError}; pub use plugin::{BuiltinTargetDescriptor, TargetPluginDescriptor, TargetPluginRegistry, TargetRequestValidator, boxed_target}; diff --git a/crates/targets/src/plugin.rs b/crates/targets/src/plugin.rs index 24743e315..4e2ce2d5f 100644 --- a/crates/targets/src/plugin.rs +++ b/crates/targets/src/plugin.rs @@ -31,7 +31,7 @@ pub enum TargetRequestValidator { Mqtt, Amqp(crate::target::TargetType), Kafka(crate::target::TargetType), - MySql, + MySql(crate::target::TargetType), Nats(crate::target::TargetType), Postgres(crate::target::TargetType), Pulsar(crate::target::TargetType), diff --git a/crates/targets/src/target/mod.rs b/crates/targets/src/target/mod.rs index 7f6581b72..d10a30d0c 100644 --- a/crates/targets/src/target/mod.rs +++ b/crates/targets/src/target/mod.rs @@ -23,7 +23,7 @@ use std::fmt::Formatter; use std::sync::Arc; use std::sync::atomic::{AtomicU64, Ordering}; use std::time::{SystemTime, UNIX_EPOCH}; -use tracing::warn; +use tracing::{debug, warn}; pub mod amqp; pub mod kafka; @@ -440,6 +440,20 @@ pub(crate) fn delete_stored_payload( } } +/// Ensures a rustls crypto provider is installed before any TLS operation. +/// +/// Multiple target modules (MySQL, Redis, Postgres, MQTT) need this because +/// each may be the first to perform a TLS handshake. Idempotent: if a +/// provider is already registered, returns immediately. +pub(crate) fn ensure_rustls_provider_installed() { + if rustls::crypto::CryptoProvider::get_default().is_some() { + return; + } + if let Err(err) = rustls::crypto::aws_lc_rs::default_provider().install_default() { + debug!("rustls provider already installed or unavailable: {err:?}"); + } +} + #[cfg(test)] mod tests { use super::*; diff --git a/crates/targets/src/target/mqtt.rs b/crates/targets/src/target/mqtt.rs index a33611ada..bf820fbe6 100644 --- a/crates/targets/src/target/mqtt.rs +++ b/crates/targets/src/target/mqtt.rs @@ -176,14 +176,6 @@ fn websocket_broker_url(broker: &Url, secure: bool) -> Result Result<(), TargetError> { if !Path::new(path).is_absolute() { return Err(TargetError::Configuration(format!("{field} must be an absolute path"))); @@ -215,7 +207,7 @@ fn build_root_store(ca_path: &str, trust_leaf_as_ca: bool) -> Result Result { - ensure_rustls_provider_installed(); + super::ensure_rustls_provider_installed(); let client_config = match tls .policy diff --git a/crates/targets/src/target/mysql.rs b/crates/targets/src/target/mysql.rs index 12f57683a..693f66d32 100644 --- a/crates/targets/src/target/mysql.rs +++ b/crates/targets/src/target/mysql.rs @@ -304,16 +304,6 @@ fn is_valid_identifier_segment(segment: &str) -> bool { true } -fn ensure_rustls_provider_installed() { - if rustls::crypto::CryptoProvider::get_default().is_some() { - return; - } - - if let Err(err) = rustls::crypto::aws_lc_rs::default_provider().install_default() { - debug!("rustls provider already installed or unavailable for mysql target: {err:?}"); - } -} - pub(crate) fn validate_table_name(table: &str) -> Result<(), TargetError> { let table = table.trim(); @@ -571,7 +561,7 @@ where .db_name(Some(dsn.database.clone())); if dsn.tls { - ensure_rustls_provider_installed(); + super::ensure_rustls_provider_installed(); let mut ssl_opts = SslOpts::default(); if !self.args.tls_ca.is_empty() { ssl_opts = ssl_opts.with_root_certs(vec![PathBuf::from(self.args.tls_ca.clone()).into()]); @@ -666,7 +656,7 @@ where conn.exec_drop(sql, (event_time.as_str(), event_data)) .await - .map_err(map_mysql_error)?; + .map_err(|err| map_mysql_error(err, "Failed to insert event"))?; self.delivery_counters.record_success(); debug!(target_id = %self.id, "MySQL event inserted"); @@ -690,16 +680,16 @@ where /// - `Server(1213|1205|1040)` → `Timeout` (deadlock/lock timeout/too /// many connections, exponential-backoff retry) /// - everything else → `Request` (permanent failure) -fn map_mysql_error(err: mysql_async::Error) -> TargetError { +pub(crate) fn map_mysql_error(err: mysql_async::Error, operation: &str) -> TargetError { match &err { mysql_async::Error::Io(_) | mysql_async::Error::Driver(_) => TargetError::NotConnected, mysql_async::Error::Server(server_err) => match server_err.code { 1213 | 1205 | 1040 => { TargetError::Timeout(format!("MySQL transient server error {}: {}", server_err.code, server_err.message)) } - _ => TargetError::Request(format!("Failed to insert event: {err}")), + _ => TargetError::Request(format!("{operation}: {err}")), }, - _ => TargetError::Request(format!("Failed to insert event: {err}")), + _ => TargetError::Request(format!("{operation}: {err}")), } } diff --git a/crates/targets/src/target/postgres.rs b/crates/targets/src/target/postgres.rs index fcd489095..7a474efb8 100644 --- a/crates/targets/src/target/postgres.rs +++ b/crates/targets/src/target/postgres.rs @@ -48,7 +48,7 @@ use std::path::{Path, PathBuf}; use std::sync::Arc; use tokio_postgres::Config; use tokio_postgres_rustls::MakeRustlsConnect; -use tracing::{debug, error, info, instrument, warn}; +use tracing::{error, info, instrument, warn}; use url::Url; use uuid::Uuid; @@ -398,14 +398,6 @@ pub fn table_probe_sql(schema: &str, table: &str) -> String { format!("SELECT 1 FROM {} LIMIT 0", qualified_table(schema, table)) } -fn ensure_rustls_provider_installed() { - if rustls::crypto::CryptoProvider::get_default().is_none() - && rustls::crypto::aws_lc_rs::default_provider().install_default().is_err() - { - debug!("rustls crypto provider was installed concurrently, skipping aws-lc-rs install"); - } -} - /// Builds a rustls `ClientConfig` for the PostgreSQL connection. /// /// When `tls_ca` is empty the OS native trust store is used via @@ -413,7 +405,7 @@ fn ensure_rustls_provider_installed() { /// When `tls_client_cert` and `tls_client_key` are both set the connection /// uses mTLS authentication; otherwise no client cert is sent. pub fn build_tls_config(args: &PostgresArgs) -> Result { - ensure_rustls_provider_installed(); + super::ensure_rustls_provider_installed(); let mut root_store = rustls::RootCertStore::empty(); diff --git a/crates/targets/src/target/redis.rs b/crates/targets/src/target/redis.rs index bcc96468a..d8fafa99a 100644 --- a/crates/targets/src/target/redis.rs +++ b/crates/targets/src/target/redis.rs @@ -290,14 +290,6 @@ fn validate_redis_tls_config(url: &Url, tls: &RedisTlsConfig) -> Result<(), Targ Ok(()) } -fn ensure_rustls_provider_installed() { - if rustls::crypto::CryptoProvider::get_default().is_none() - && rustls::crypto::aws_lc_rs::default_provider().install_default().is_err() - { - debug!("rustls crypto provider was installed concurrently, skipping aws-lc-rs install"); - } -} - pub struct RedisTarget where E: Send + Sync + 'static + Clone + Serialize + DeserializeOwned, @@ -656,7 +648,7 @@ pub(crate) fn build_redis_client(args: &RedisArgs) -> Result Arc::new(validate_mysql_request_entry), + TargetRequestValidator::MySql(target_type) => { + Arc::new(move |kv_map, default_queue_dir| validate_mysql_request_entry(kv_map, default_queue_dir, target_type)) + } TargetRequestValidator::Nats(target_type) => { if matches!(TargetDomain::from(target_type), TargetDomain::Audit) { Arc::new(validate_audit_nats_request_entry) @@ -622,10 +624,11 @@ fn validate_audit_pulsar_request_entry( fn validate_mysql_request_entry( kv_map: &HashMap, default_queue_dir: &str, + target_type: TargetType, ) -> futures::future::BoxFuture<'static, S3Result<()>> { let kv_map = kv_map.clone(); let default_queue_dir = default_queue_dir.to_string(); - Box::pin(async move { validate_mysql_request(&kv_map, &default_queue_dir).await }) + Box::pin(async move { validate_mysql_request(&kv_map, &default_queue_dir, target_type).await }) } fn validate_notify_postgres_request_entry( @@ -722,14 +725,21 @@ async fn validate_pulsar_request( }) } -async fn validate_mysql_request(kv_map: &HashMap, default_queue_dir: &str) -> S3Result<()> { +async fn validate_mysql_request( + kv_map: &HashMap, + default_queue_dir: &str, + target_type: TargetType, +) -> S3Result<()> { if let Some(queue_dir) = kv_map.get(MYSQL_QUEUE_DIR) { validate_queue_dir(queue_dir.as_str()).await?; } - validate_mysql_config(&to_kvs(kv_map), default_queue_dir).map_err(|e| s3_error!(InvalidArgument, "{}", e))?; - - Ok(()) + let args = + build_mysql_args(&to_kvs(kv_map), default_queue_dir, target_type).map_err(|e| s3_error!(InvalidArgument, "{}", e))?; + check_mysql_server_available(&args).await.map_err(|e| match e { + TargetError::Configuration(_) => s3_error!(InvalidArgument, "{}", e), + _ => s3_error!(InvalidArgument, "MySQL server check failed: {}", e), + }) } async fn validate_postgres_request(