From c3586d9af6f20ad0f16cc4194a04ae9550cfc0c4 Mon Sep 17 00:00:00 2001 From: mdecimus <11444311+mdecimus@users.noreply.github.com> Date: Wed, 25 Feb 2026 14:33:52 +0100 Subject: [PATCH] SQLite: Fix thread pool exhaustion --- crates/store/src/backend/sqlite/blob.rs | 9 ++++++--- crates/store/src/backend/sqlite/lookup.rs | 3 ++- crates/store/src/backend/sqlite/main.rs | 23 +++++++++++++++-------- crates/store/src/backend/sqlite/read.rs | 12 +++++++----- crates/store/src/backend/sqlite/write.rs | 23 ++++++++--------------- 5 files changed, 38 insertions(+), 32 deletions(-) diff --git a/crates/store/src/backend/sqlite/blob.rs b/crates/store/src/backend/sqlite/blob.rs index 623496a9..29f05789 100644 --- a/crates/store/src/backend/sqlite/blob.rs +++ b/crates/store/src/backend/sqlite/blob.rs @@ -16,8 +16,9 @@ impl SqliteStore { key: &[u8], range: Range, ) -> trc::Result>> { - let conn = self.conn_pool.get().map_err(into_error)?; + let manager = self.conn_pool.clone(); self.spawn_worker(move || { + let conn = manager.get().map_err(into_error)?; let mut result = conn .prepare_cached("SELECT v FROM t WHERE k = ?") .map_err(into_error)?; @@ -42,8 +43,9 @@ impl SqliteStore { } pub(crate) async fn put_blob(&self, key: &[u8], data: &[u8]) -> trc::Result<()> { - let conn = self.conn_pool.get().map_err(into_error)?; + let manager = self.conn_pool.clone(); self.spawn_worker(move || { + let conn = manager.get().map_err(into_error)?; conn.prepare_cached("INSERT OR REPLACE INTO t (k, v) VALUES (?, ?)") .map_err(into_error)? .execute([key, data]) @@ -54,8 +56,9 @@ impl SqliteStore { } pub(crate) async fn delete_blob(&self, key: &[u8]) -> trc::Result { - let conn = self.conn_pool.get().map_err(into_error)?; + let manager = self.conn_pool.clone(); self.spawn_worker(move || { + let conn = manager.get().map_err(into_error)?; conn.prepare_cached("DELETE FROM t WHERE k = ?") .map_err(into_error)? .execute([key]) diff --git a/crates/store/src/backend/sqlite/lookup.rs b/crates/store/src/backend/sqlite/lookup.rs index 937dfff0..8dfbbd2f 100644 --- a/crates/store/src/backend/sqlite/lookup.rs +++ b/crates/store/src/backend/sqlite/lookup.rs @@ -16,8 +16,9 @@ impl SqliteStore { query: &str, params_: &[Value<'_>], ) -> trc::Result { - let conn = self.conn_pool.get().map_err(into_error)?; + let manager = self.conn_pool.clone(); self.spawn_worker(move || { + let conn = manager.get().map_err(into_error)?; let mut s = conn.prepare_cached(query).map_err(into_error)?; let params = params_ .iter() diff --git a/crates/store/src/backend/sqlite/main.rs b/crates/store/src/backend/sqlite/main.rs index c40dad0e..38b9eeb1 100644 --- a/crates/store/src/backend/sqlite/main.rs +++ b/crates/store/src/backend/sqlite/main.rs @@ -79,6 +79,7 @@ impl SqliteStore { SUBSPACE_TELEMETRY_SPAN, SUBSPACE_TELEMETRY_METRIC, SUBSPACE_SEARCH_INDEX, + SUBSPACE_DIRECTORY, ] { let table = char::from(table); conn.execute( @@ -93,16 +94,22 @@ impl SqliteStore { .map_err(into_error)?; } - let table = char::from(SUBSPACE_INDEXES); - conn.execute( - &format!( - "CREATE TABLE IF NOT EXISTS {table} ( + for table in [ + SUBSPACE_INDEXES, + SUBSPACE_REGISTRY_IDX, + SUBSPACE_REGISTRY_IDX_GLOBAL, + ] { + let table = char::from(table); + conn.execute( + &format!( + "CREATE TABLE IF NOT EXISTS {table} ( k BLOB PRIMARY KEY )" - ), - [], - ) - .map_err(into_error)?; + ), + [], + ) + .map_err(into_error)?; + } for table in [SUBSPACE_COUNTER, SUBSPACE_QUOTA, SUBSPACE_IN_MEMORY_COUNTER] { conn.execute( diff --git a/crates/store/src/backend/sqlite/read.rs b/crates/store/src/backend/sqlite/read.rs index 906bb3dc..b9af98dd 100644 --- a/crates/store/src/backend/sqlite/read.rs +++ b/crates/store/src/backend/sqlite/read.rs @@ -13,8 +13,9 @@ impl SqliteStore { where U: Deserialize + 'static, { - let conn = self.conn_pool.get().map_err(into_error)?; + let manager = self.conn_pool.clone(); self.spawn_worker(move || { + let conn = manager.get().map_err(into_error)?; let mut result = conn .prepare_cached(&format!( "SELECT v FROM {} WHERE k = ?", @@ -24,7 +25,7 @@ impl SqliteStore { let key = key.serialize(0); result .query_row([&key], |row| { - U::deserialize(row.get_ref(0)?.as_bytes()?) + U::deserialize_with_key(&key, row.get_ref(0)?.as_bytes()?) .map_err(|err| rusqlite::Error::ToSqlConversionFailure(err.into())) }) .optional() @@ -38,9 +39,9 @@ impl SqliteStore { params: IterateParams, mut cb: impl for<'x> FnMut(&'x [u8], &'x [u8]) -> trc::Result + Sync + Send, ) -> trc::Result<()> { - let conn = self.conn_pool.get().map_err(into_error)?; - + let manager = self.conn_pool.clone(); self.spawn_worker(move || { + let conn = manager.get().map_err(into_error)?; let table = char::from(params.begin.subspace()); let begin = params.begin.serialize(0); let end = params.end.serialize(0); @@ -113,8 +114,9 @@ impl SqliteStore { let key = key.into(); let table = char::from(key.subspace()); let key = key.serialize(0); - let conn = self.conn_pool.get().map_err(into_error)?; + let manager = self.conn_pool.clone(); self.spawn_worker(move || { + let conn = manager.get().map_err(into_error)?; match conn .prepare_cached(&format!("SELECT v FROM {table} WHERE k = ?")) .map_err(into_error)? diff --git a/crates/store/src/backend/sqlite/write.rs b/crates/store/src/backend/sqlite/write.rs index e96d0a90..c4ba38ca 100644 --- a/crates/store/src/backend/sqlite/write.rs +++ b/crates/store/src/backend/sqlite/write.rs @@ -14,12 +14,10 @@ use trc::AddContext; impl SqliteStore { pub(crate) async fn write(&self, batch: Batch<'_>) -> trc::Result { - let mut conn = self - .conn_pool - .get() - .map_err(into_error) - .caused_by(trc::location!())?; + let manager = self.conn_pool.clone(); self.spawn_worker(move || { + let mut conn = manager.get().map_err(into_error)?; + let mut account_id = u32::MAX; let mut collection = u8::MAX; let mut document_id = u32::MAX; @@ -271,12 +269,9 @@ impl SqliteStore { } pub(crate) async fn purge_store(&self) -> trc::Result<()> { - let conn = self - .conn_pool - .get() - .map_err(into_error) - .caused_by(trc::location!())?; + let manager = self.conn_pool.clone(); self.spawn_worker(move || { + let conn = manager.get().map_err(into_error)?; for subspace in [SUBSPACE_QUOTA, SUBSPACE_COUNTER, SUBSPACE_IN_MEMORY_COUNTER] { conn.prepare_cached(&format!("DELETE FROM {} WHERE v = 0", char::from(subspace),)) .map_err(into_error) @@ -292,12 +287,10 @@ impl SqliteStore { } pub(crate) async fn delete_range(&self, from: impl Key, to: impl Key) -> trc::Result<()> { - let conn = self - .conn_pool - .get() - .map_err(into_error) - .caused_by(trc::location!())?; + let manager = self.conn_pool.clone(); self.spawn_worker(move || { + let conn = manager.get().map_err(into_error)?; + conn.prepare_cached(&format!( "DELETE FROM {} WHERE k >= ? AND k < ?", char::from(from.subspace()),