From ff26e90177579409e58045a2b6b58c9a14cce6df Mon Sep 17 00:00:00 2001 From: t8y2 <1156263951@qq.com> Date: Sun, 31 May 2026 00:49:45 +0800 Subject: [PATCH] fix(duckdb): reuse connection for query sessions --- crates/dbx-core/src/connection.rs | 28 ++++++++++++++++++++++------ 1 file changed, 22 insertions(+), 6 deletions(-) diff --git a/crates/dbx-core/src/connection.rs b/crates/dbx-core/src/connection.rs index 63364d410..21f0e23ed 100644 --- a/crates/dbx-core/src/connection.rs +++ b/crates/dbx-core/src/connection.rs @@ -155,7 +155,7 @@ impl AppState { }; let base_pool_key = base_pool_key_for(db_type, connection_id, database, false); - let pool_key = session_scoped_pool_key(base_pool_key, client_session_id); + let pool_key = session_scoped_pool_key_for(db_type, base_pool_key, client_session_id); let conns = self.connections.read().await; if conns.contains_key(&pool_key) { @@ -491,7 +491,7 @@ impl AppState { configs.get(connection_id).map(|c| c.db_type) }; let base_pool_key = base_pool_key_for(db_type, connection_id, database, true); - let pool_key = session_scoped_pool_key(base_pool_key, client_session_id); + let pool_key = session_scoped_pool_key_for(db_type, base_pool_key, client_session_id); if self.uses_forwarded_transport(connection_id).await { self.remove_connection_pools(connection_id).await; self.reset_connection_transport(connection_id).await; @@ -519,7 +519,10 @@ impl AppState { configs.get(connection_id).map(|c| c.db_type) }; let base_pool_key = base_pool_key_for(db_type, connection_id, database, false); - let pool_key = session_scoped_pool_key(base_pool_key, Some(&session)); + let pool_key = session_scoped_pool_key_for(db_type, base_pool_key.clone(), Some(&session)); + if pool_key == base_pool_key { + return Ok(false); + } let removed = self.connections.write().await.remove(&pool_key); if let Some(pool) = removed { close_pool_kind(pool).await; @@ -704,6 +707,17 @@ fn session_scoped_pool_key(base_pool_key: String, client_session_id: Option<&str .unwrap_or(base_pool_key) } +fn session_scoped_pool_key_for( + db_type: Option, + base_pool_key: String, + client_session_id: Option<&str>, +) -> String { + if matches!(db_type, Some(DatabaseType::DuckDb)) { + return base_pool_key; + } + session_scoped_pool_key(base_pool_key, client_session_id) +} + fn clone_pool_kind(pool: &PoolKind) -> PoolKind { match pool { PoolKind::Mysql(p, mode) => PoolKind::Mysql(p.clone(), *mode), @@ -1479,7 +1493,7 @@ mod tests { } #[tokio::test] - async fn close_client_session_pool_removes_session_scoped_duckdb_pool() { + async fn duckdb_client_session_reuses_base_pool_to_avoid_file_locks() { let (state, dir) = test_app_state().await; let db_path = dir.join("session.duckdb"); let mut config = mysql_config(None); @@ -1490,12 +1504,14 @@ mod tests { config.port = 0; state.configs.write().await.insert(config.id.clone(), config.clone()); + let base_pool_key = state.get_or_create_pool("duckdb-conn", None).await.unwrap(); let pool_key = state.get_or_create_pool_for_session("duckdb-conn", Some("main"), Some("tab-1")).await.unwrap(); - assert_eq!(pool_key, "duckdb-conn:session:tab-1"); + assert_eq!(pool_key, base_pool_key); - assert!(state.close_client_session_pool("duckdb-conn", Some("main"), "tab-1").await.unwrap()); + assert!(!state.close_client_session_pool("duckdb-conn", Some("main"), "tab-1").await.unwrap()); let conns = state.connections.read().await; + assert!(conns.contains_key("duckdb-conn")); assert!(!conns.contains_key("duckdb-conn:session:tab-1")); let _ = std::fs::remove_dir_all(dir);