From 0fe25f0eba2ad7657c873fb1dff3addc15a94a4d Mon Sep 17 00:00:00 2001 From: lanlan-wzz <54491755+lanlan-wzz@users.noreply.github.com> Date: Wed, 29 Jul 2026 00:08:04 +0800 Subject: [PATCH] fix(mysql): retry legacy EOF mode for incompatible proxies --- Cargo.lock | 2 +- Cargo.toml | 2 +- crates/dbx-core/src/db/mysql.rs | 172 +++++++++++++++++++++++++------- 3 files changed, 138 insertions(+), 38 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 010f3b3ba..a383ed13a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4755,7 +4755,7 @@ dependencies = [ [[package]] name = "mysql_async" version = "0.37.0" -source = "git+https://github.com/t8y2/mysql_async.git?rev=5c1186d05e62f712e717632a5b66dd030466e9b3#5c1186d05e62f712e717632a5b66dd030466e9b3" +source = "git+https://github.com/t8y2/mysql_async.git?rev=2be6e392eb9b06d20dcd2d8ed8eae748d413c9ec#2be6e392eb9b06d20dcd2d8ed8eae748d413c9ec" dependencies = [ "bytes", "crossbeam-queue", diff --git a/Cargo.toml b/Cargo.toml index 716cfc927..e619d0635 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -12,7 +12,7 @@ postgres-types = { git = "https://github.com/t8y2/tokio-postgres-gaussdb.git", b postgres-protocol = { git = "https://github.com/t8y2/tokio-postgres-gaussdb.git", branch = "master" } # Temporary compatibility patch for MySQL accounts using the deprecated # sha256_password auth plugin. Remove after upstream support lands. -mysql_async = { git = "https://github.com/t8y2/mysql_async.git", rev = "5c1186d05e62f712e717632a5b66dd030466e9b3" } +mysql_async = { git = "https://github.com/t8y2/mysql_async.git", rev = "2be6e392eb9b06d20dcd2d8ed8eae748d413c9ec" } [profile.release] panic = "abort" diff --git a/crates/dbx-core/src/db/mysql.rs b/crates/dbx-core/src/db/mysql.rs index 1050cc639..e8590c400 100644 --- a/crates/dbx-core/src/db/mysql.rs +++ b/crates/dbx-core/src/db/mysql.rs @@ -637,6 +637,72 @@ async fn connect_with_ca_cert_pool_limit_idle_setup_database_with_mode( setup_mode: MySqlSetupMode, ) -> Result { let timeout = super::parse_connect_timeout_with_fallback(url, fallback_timeout); + let mut retry_url = url.to_string(); + let mut retry_ca_cert_path = ca_cert_path; + let mut result = connect_pool_attempt( + url, + ca_cert_path, + timeout, + max_connections, + idle_timeout_secs, + setup_database, + extra_setup_queries, + setup_mode, + MySqlEofMode::Deprecate, + ) + .await; + + if result.as_ref().err().is_some_and(|error| mysql_error_should_retry_without_ssl(error)) { + if let Some(fallback_url) = ssl_fallback_url(url) { + log::info!("SSL handshake failed, retrying with ssl-mode=disabled"); + retry_url = fallback_url; + retry_ca_cert_path = None; + result = connect_pool_attempt( + &retry_url, + None, + timeout, + max_connections, + idle_timeout_secs, + setup_database, + extra_setup_queries, + setup_mode, + MySqlEofMode::Deprecate, + ) + .await; + } + } + + if result.as_ref().err().is_some_and(|error| mysql_error_should_retry_with_legacy_eof(error)) { + log::info!("MySQL proxy returned legacy EOF packets; retrying with CLIENT_DEPRECATE_EOF disabled"); + return connect_pool_attempt( + &retry_url, + retry_ca_cert_path, + timeout, + max_connections, + idle_timeout_secs, + setup_database, + extra_setup_queries, + setup_mode, + MySqlEofMode::Legacy, + ) + .await; + } + + result +} + +#[allow(clippy::too_many_arguments)] +async fn connect_pool_attempt( + url: &str, + ca_cert_path: Option<&str>, + timeout: Duration, + max_connections: usize, + idle_timeout_secs: Option, + setup_database: Option<&str>, + extra_setup_queries: &[String], + setup_mode: MySqlSetupMode, + eof_mode: MySqlEofMode, +) -> Result { let pool = create_pool( url, ca_cert_path, @@ -645,8 +711,9 @@ async fn connect_with_ca_cert_pool_limit_idle_setup_database_with_mode( setup_database, extra_setup_queries, setup_mode, + eof_mode, )?; - let result = verify_pool_connection_with_setup_fallback( + verify_pool_connection_with_setup_fallback( pool, timeout, url, @@ -656,39 +723,9 @@ async fn connect_with_ca_cert_pool_limit_idle_setup_database_with_mode( setup_database, extra_setup_queries, setup_mode, + eof_mode, ) - .await; - - if let Err(ref e) = result { - if mysql_error_should_retry_without_ssl(e) { - if let Some(fallback_url) = ssl_fallback_url(url) { - log::info!("SSL handshake failed, retrying with ssl-mode=disabled"); - let fallback_pool = create_pool( - &fallback_url, - None, - max_connections, - idle_timeout_secs, - setup_database, - extra_setup_queries, - setup_mode, - )?; - return verify_pool_connection_with_setup_fallback( - fallback_pool, - timeout, - &fallback_url, - None, - max_connections, - idle_timeout_secs, - setup_database, - extra_setup_queries, - setup_mode, - ) - .await; - } - } - } - - result + .await } #[derive(Debug, Default, Clone, PartialEq, Eq)] @@ -703,6 +740,18 @@ enum MySqlSetupMode { Compatible, } +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum MySqlEofMode { + Deprecate, + Legacy, +} + +impl MySqlEofMode { + fn deprecate_eof(self) -> bool { + self == Self::Deprecate + } +} + const MYSQL_GROUP_CONCAT_MAX_LEN: u64 = 1_048_576; impl MySqlSetupMode { @@ -724,6 +773,7 @@ async fn verify_pool_connection_with_setup_fallback( setup_database: Option<&str>, extra_setup_queries: &[String], setup_mode: MySqlSetupMode, + eof_mode: MySqlEofMode, ) -> Result { match verify_pool_connection(&pool, timeout).await { Ok(()) => Ok(pool), @@ -742,6 +792,7 @@ async fn verify_pool_connection_with_setup_fallback( setup_database, extra_setup_queries, fallback_mode, + eof_mode, )?; verify_pool_connection(&fallback_pool, timeout).await.map(|_| fallback_pool) } @@ -783,6 +834,7 @@ fn create_pool( setup_database: Option<&str>, extra_setup_queries: &[String], setup_mode: MySqlSetupMode, + eof_mode: MySqlEofMode, ) -> Result { let tls_url = mysql_tls_url(url)?; let local_infile_paths = mysql_local_infile_paths(&tls_url.url); @@ -816,6 +868,7 @@ fn create_pool( .prefer_socket(false) .pool_opts(Some(pool_opts)) .tcp_keepalive(Some(Duration::from_millis(u64::from(MYSQL_TCP_KEEPALIVE_MS)))) + .deprecate_eof(eof_mode.deprecate_eof()) .setup(setup_queries); if let Some(ssl_opts) = mysql_ssl_opts(base_ssl_opts, url, ca_cert_path, &tls_url.files)? { builder = builder.ssl_opts(ssl_opts); @@ -1323,6 +1376,10 @@ fn mysql_error_should_retry_without_ssl(error: &str) -> bool { || (error.contains("client asked for ssl") && error.contains("server does not have this capability")) } +fn mysql_error_should_retry_with_legacy_eof(error: &str) -> bool { + error.to_ascii_lowercase().contains("packets out of sync") +} + fn mysql_error_should_retry_with_text_protocol(error: &str) -> bool { let lower = error.to_ascii_lowercase(); (lower.contains("1105") && lower.contains("hy000")) @@ -1658,9 +1715,36 @@ pub async fn connect_bare_with_pool_limit_and_setup_database( extra_setup_queries: &[String], ) -> Result { let timeout = super::parse_connect_timeout_with_fallback(url, fallback_timeout); - let pool = - create_pool(url, None, max_connections, None, setup_database, extra_setup_queries, MySqlSetupMode::Compatible)?; - verify_pool_connection(&pool, timeout).await.map(|_| pool) + let result = connect_pool_attempt( + url, + None, + timeout, + max_connections, + None, + setup_database, + extra_setup_queries, + MySqlSetupMode::Compatible, + MySqlEofMode::Deprecate, + ) + .await; + if result.as_ref().err().is_some_and(|error| mysql_error_should_retry_with_legacy_eof(error)) { + log::info!( + "MySQL proxy returned legacy EOF packets; retrying bare connection with CLIENT_DEPRECATE_EOF disabled" + ); + return connect_pool_attempt( + url, + None, + timeout, + max_connections, + None, + setup_database, + extra_setup_queries, + MySqlSetupMode::Compatible, + MySqlEofMode::Legacy, + ) + .await; + } + result } pub async fn list_databases(pool: &MySqlPool) -> Result, String> { @@ -4951,6 +5035,22 @@ mod tests { let error = "MySQL connection failed: Input/output error: Input/output error: packet out of order"; assert!(mysql_error_should_retry_without_ssl(error)); + assert!(!mysql_error_should_retry_with_legacy_eof(error)); + } + + #[test] + fn mysql_packets_out_of_sync_retries_with_legacy_eof() { + let error = "MySQL connection failed: Input/output error: Input/output error: Packets out of sync"; + + assert!(mysql_error_should_retry_with_legacy_eof(error)); + assert!(!mysql_error_should_retry_without_ssl(error)); + } + + #[test] + fn mysql_async_builder_can_disable_deprecated_eof_protocol() { + let opts = mysql_async::Opts::from(mysql_async::OptsBuilder::default().deprecate_eof(false)); + + assert!(!opts.deprecate_eof()); } #[test]