This is an automated email from the ASF dual-hosted git repository.

erickguan pushed a commit to branch fix-stat-impl
in repository https://gitbox.apache.org/repos/asf/opendal.git


The following commit(s) were added to refs/heads/fix-stat-impl by this push:
     new 2b4007c1f refactor: cache sql service queries
2b4007c1f is described below

commit 2b4007c1f97d59eb3fef3a02e2b606c88e2d7f4e
Author: Erick Guan <[email protected]>
AuthorDate: Sun Jul 5 13:02:44 2026 +0800

    refactor: cache sql service queries
---
 core/services/mysql/src/backend.rs      |   6 +-
 core/services/mysql/src/core.rs         | 158 +++++++++++++++++-------
 core/services/postgresql/src/backend.rs |   6 +-
 core/services/postgresql/src/core.rs    | 141 +++++++++++++++------
 core/services/sqlite/src/backend.rs     |  20 +--
 core/services/sqlite/src/core.rs        | 210 +++++++++++++++++++++++---------
 6 files changed, 388 insertions(+), 153 deletions(-)

diff --git a/core/services/mysql/src/backend.rs 
b/core/services/mysql/src/backend.rs
index da2b672f2..a9634b583 100644
--- a/core/services/mysql/src/backend.rs
+++ b/core/services/mysql/src/backend.rs
@@ -139,13 +139,13 @@ impl Builder for MysqlBuilder {
 
         let root = normalize_root(self.config.root.unwrap_or_else(|| 
"/".to_string()).as_str());
 
-        Ok(MysqlBackend::new(MysqlCore {
-            pool: OnceCell::new(),
+        Ok(MysqlBackend::new(MysqlCore::new(
+            OnceCell::new(),
             config,
             table,
             key_field,
             value_field,
-        })
+        ))
         .with_normalized_root(root))
     }
 }
diff --git a/core/services/mysql/src/core.rs b/core/services/mysql/src/core.rs
index 961c29807..4b2199586 100644
--- a/core/services/mysql/src/core.rs
+++ b/core/services/mysql/src/core.rs
@@ -15,6 +15,8 @@
 // specific language governing permissions and limitations
 // under the License.
 
+use std::sync::OnceLock;
+
 use mea::once::OnceCell;
 use sqlx::MySqlPool;
 use sqlx::mysql::MySqlConnectOptions;
@@ -29,9 +31,37 @@ pub struct MysqlCore {
     pub table: String,
     pub key_field: String,
     pub value_field: String,
+
+    queries: MysqlQueries,
+}
+
+#[derive(Clone, Debug, Default)]
+struct MysqlQueries {
+    get: OnceLock<String>,
+    get_length: OnceLock<String>,
+    set: OnceLock<String>,
+    delete: OnceLock<String>,
+    list: OnceLock<String>,
 }
 
 impl MysqlCore {
+    pub(crate) fn new(
+        pool: OnceCell<MySqlPool>,
+        config: MySqlConnectOptions,
+        table: String,
+        key_field: String,
+        value_field: String,
+    ) -> Self {
+        Self {
+            pool,
+            config,
+            table,
+            key_field,
+            value_field,
+            queries: MysqlQueries::default(),
+        }
+    }
+
     async fn get_client(&self) -> Result<&MySqlPool> {
         self.pool
             .get_or_try_init(|| async {
@@ -42,6 +72,71 @@ impl MysqlCore {
             .await
     }
 
+    fn query_get(&self) -> &str {
+        self.queries
+            .get
+            .get_or_init(|| {
+                format!(
+                    "SELECT `{}` FROM `{}` WHERE `{}` = ? LIMIT 1",
+                    self.value_field, self.table, self.key_field
+                )
+            })
+            .as_str()
+    }
+
+    fn query_get_length(&self) -> &str {
+        self.queries
+            .get_length
+            .get_or_init(|| {
+                format!(
+                    "SELECT OCTET_LENGTH(`{}`) FROM `{}` WHERE `{}` = ? LIMIT 
1",
+                    self.value_field, self.table, self.key_field
+                )
+            })
+            .as_str()
+    }
+
+    fn query_set(&self) -> &str {
+        self.queries
+            .set
+            .get_or_init(|| {
+                format!(
+                    r#"INSERT INTO `{}` (`{}`, `{}`) VALUES (?, ?)
+            ON DUPLICATE KEY UPDATE `{}` = VALUES(`{}`)"#,
+                    self.table,
+                    self.key_field,
+                    self.value_field,
+                    self.value_field,
+                    self.value_field
+                )
+            })
+            .as_str()
+    }
+
+    fn query_delete(&self) -> &str {
+        self.queries
+            .delete
+            .get_or_init(|| {
+                format!(
+                    "DELETE FROM `{}` WHERE `{}` = ?",
+                    self.table, self.key_field
+                )
+            })
+            .as_str()
+    }
+
+    fn query_list(&self) -> &str {
+        self.queries
+            .list
+            .get_or_init(|| {
+                format!(
+                    "SELECT `{}` FROM `{}` WHERE `{}` LIKE ? ORDER BY `{}`",
+                    self.key_field, self.table, self.key_field, self.key_field
+                )
+            })
+            .as_str()
+    }
+
     pub async fn get(&self, path: &str) -> Result<Option<Buffer>> {
         let pool = self.get_client().await?;
 
@@ -55,14 +150,11 @@ impl MysqlCore {
         // - formatted, quoted identifiers for trusted table and field 
configuration
         // - bind parameters for values to avoid malformed SQL and SQL 
injection
         // to ensure correctness.
-        let value: Option<Vec<u8>> = sqlx::query_scalar(&format!(
-            "SELECT `{}` FROM `{}` WHERE `{}` = ? LIMIT 1",
-            self.value_field, self.table, self.key_field
-        ))
-        .bind(path)
-        .fetch_optional(pool)
-        .await
-        .map_err(parse_mysql_error)?;
+        let value: Option<Vec<u8>> = sqlx::query_scalar(self.query_get())
+            .bind(path)
+            .fetch_optional(pool)
+            .await
+            .map_err(parse_mysql_error)?;
 
         Ok(value.map(Buffer::from))
     }
@@ -70,14 +162,11 @@ impl MysqlCore {
     pub async fn get_length(&self, path: &str) -> Result<Option<usize>> {
         let pool = self.get_client().await?;
 
-        let value: Option<i64> = sqlx::query_scalar(&format!(
-            "SELECT OCTET_LENGTH(`{}`) FROM `{}` WHERE `{}` = ? LIMIT 1",
-            self.value_field, self.table, self.key_field
-        ))
-        .bind(path)
-        .fetch_optional(pool)
-        .await
-        .map_err(parse_mysql_error)?;
+        let value: Option<i64> = sqlx::query_scalar(self.query_get_length())
+            .bind(path)
+            .fetch_optional(pool)
+            .await
+            .map_err(parse_mysql_error)?;
 
         value
             .map(|v| {
@@ -92,16 +181,12 @@ impl MysqlCore {
     pub async fn set(&self, path: &str, value: Buffer) -> Result<()> {
         let pool = self.get_client().await?;
 
-        sqlx::query(&format!(
-            r#"INSERT INTO `{}` (`{}`, `{}`) VALUES (?, ?)
-            ON DUPLICATE KEY UPDATE `{}` = VALUES(`{}`)"#,
-            self.table, self.key_field, self.value_field, self.value_field, 
self.value_field
-        ))
-        .bind(path)
-        .bind(value.to_vec())
-        .execute(pool)
-        .await
-        .map_err(parse_mysql_error)?;
+        sqlx::query(self.query_set())
+            .bind(path)
+            .bind(value.to_vec())
+            .execute(pool)
+            .await
+            .map_err(parse_mysql_error)?;
 
         Ok(())
     }
@@ -109,14 +194,11 @@ impl MysqlCore {
     pub async fn delete(&self, path: &str) -> Result<()> {
         let pool = self.get_client().await?;
 
-        sqlx::query(&format!(
-            "DELETE FROM `{}` WHERE `{}` = ?",
-            self.table, self.key_field
-        ))
-        .bind(path)
-        .execute(pool)
-        .await
-        .map_err(parse_mysql_error)?;
+        sqlx::query(self.query_delete())
+            .bind(path)
+            .execute(pool)
+            .await
+            .map_err(parse_mysql_error)?;
 
         Ok(())
     }
@@ -124,14 +206,8 @@ impl MysqlCore {
     pub async fn list(&self, path: &str) -> Result<Vec<String>> {
         let pool = self.get_client().await?;
 
-        let mut sql = format!(
-            "SELECT `{}` FROM `{}` WHERE `{}` LIKE ?",
-            self.key_field, self.table, self.key_field
-        );
-        sql.push_str(&format!(" ORDER BY `{}`", self.key_field));
-
         let escaped = escape_like(path);
-        sqlx::query_scalar(&sql)
+        sqlx::query_scalar(self.query_list())
             .bind(format!("{escaped}%"))
             .fetch_all(pool)
             .await
diff --git a/core/services/postgresql/src/backend.rs 
b/core/services/postgresql/src/backend.rs
index ebd33a007..de73412e6 100644
--- a/core/services/postgresql/src/backend.rs
+++ b/core/services/postgresql/src/backend.rs
@@ -133,13 +133,13 @@ impl Builder for PostgresqlBuilder {
 
         let root = normalize_root(self.config.root.unwrap_or_else(|| 
"/".to_string()).as_str());
 
-        Ok(PostgresqlBackend::new(PostgresqlCore {
-            pool: OnceCell::new(),
+        Ok(PostgresqlBackend::new(PostgresqlCore::new(
+            OnceCell::new(),
             config,
             table,
             key_field,
             value_field,
-        })
+        ))
         .with_normalized_root(root))
     }
 }
diff --git a/core/services/postgresql/src/core.rs 
b/core/services/postgresql/src/core.rs
index ff6107ab1..5a8acc1c7 100644
--- a/core/services/postgresql/src/core.rs
+++ b/core/services/postgresql/src/core.rs
@@ -15,6 +15,8 @@
 // specific language governing permissions and limitations
 // under the License.
 
+use std::sync::OnceLock;
+
 use mea::once::OnceCell;
 use sqlx::PgPool;
 use sqlx::postgres::PgConnectOptions;
@@ -29,9 +31,36 @@ pub struct PostgresqlCore {
     pub table: String,
     pub key_field: String,
     pub value_field: String,
+
+    queries: PostgresqlQueries,
+}
+
+#[derive(Clone, Debug, Default)]
+struct PostgresqlQueries {
+    get: OnceLock<String>,
+    get_length: OnceLock<String>,
+    set: OnceLock<String>,
+    delete: OnceLock<String>,
 }
 
 impl PostgresqlCore {
+    pub(crate) fn new(
+        pool: OnceCell<PgPool>,
+        config: PgConnectOptions,
+        table: String,
+        key_field: String,
+        value_field: String,
+    ) -> Self {
+        Self {
+            pool,
+            config,
+            table,
+            key_field,
+            value_field,
+            queries: PostgresqlQueries::default(),
+        }
+    }
+
     async fn get_client(&self) -> Result<&PgPool> {
         self.pool
             .get_or_try_init(|| async {
@@ -42,6 +71,59 @@ impl PostgresqlCore {
             .await
     }
 
+    fn query_get(&self) -> &str {
+        self.queries
+            .get
+            .get_or_init(|| {
+                format!(
+                    r#"SELECT "{}" FROM "{}" WHERE "{}" = $1 LIMIT 1"#,
+                    self.value_field, self.table, self.key_field
+                )
+            })
+            .as_str()
+    }
+
+    fn query_get_length(&self) -> &str {
+        self.queries
+            .get_length
+            .get_or_init(|| {
+                format!(
+                    r#"SELECT OCTET_LENGTH("{}")::BIGINT FROM "{}" WHERE "{}" 
= $1 LIMIT 1"#,
+                    self.value_field, self.table, self.key_field
+                )
+            })
+            .as_str()
+    }
+
+    fn query_set(&self) -> &str {
+        self.queries
+            .set
+            .get_or_init(|| {
+                let table = &self.table;
+                let key_field = &self.key_field;
+                let value_field = &self.value_field;
+                format!(
+                    r#"INSERT INTO "{table}" ("{key_field}", "{value_field}")
+                VALUES ($1, $2)
+                ON CONFLICT ("{key_field}")
+                    DO UPDATE SET "{value_field}" = EXCLUDED."{value_field}""#,
+                )
+            })
+            .as_str()
+    }
+
+    fn query_delete(&self) -> &str {
+        self.queries
+            .delete
+            .get_or_init(|| {
+                format!(
+                    r#"DELETE FROM "{}" WHERE "{}" = $1"#,
+                    self.table, self.key_field
+                )
+            })
+            .as_str()
+    }
+
     pub async fn get(&self, path: &str) -> Result<Option<Buffer>> {
         let pool = self.get_client().await?;
 
@@ -53,14 +135,11 @@ impl PostgresqlCore {
         // - formatted, quoted identifiers for trusted table and field 
configuration
         // - bind parameters for values to avoid malformed SQL and SQL 
injection
         // to ensure correctness.
-        let value: Option<Vec<u8>> = sqlx::query_scalar(&format!(
-            r#"SELECT "{}" FROM "{}" WHERE "{}" = $1 LIMIT 1"#,
-            self.value_field, self.table, self.key_field
-        ))
-        .bind(path)
-        .fetch_optional(pool)
-        .await
-        .map_err(parse_postgres_error)?;
+        let value: Option<Vec<u8>> = sqlx::query_scalar(self.query_get())
+            .bind(path)
+            .fetch_optional(pool)
+            .await
+            .map_err(parse_postgres_error)?;
 
         Ok(value.map(Buffer::from))
     }
@@ -68,14 +147,11 @@ impl PostgresqlCore {
     pub async fn get_length(&self, path: &str) -> Result<Option<usize>> {
         let pool = self.get_client().await?;
 
-        let value: Option<i64> = sqlx::query_scalar(&format!(
-            r#"SELECT OCTET_LENGTH("{}")::BIGINT FROM "{}" WHERE "{}" = $1 
LIMIT 1"#,
-            self.value_field, self.table, self.key_field
-        ))
-        .bind(path)
-        .fetch_optional(pool)
-        .await
-        .map_err(parse_postgres_error)?;
+        let value: Option<i64> = sqlx::query_scalar(self.query_get_length())
+            .bind(path)
+            .fetch_optional(pool)
+            .await
+            .map_err(parse_postgres_error)?;
 
         value
             .map(|v| {
@@ -90,20 +166,12 @@ impl PostgresqlCore {
     pub async fn set(&self, path: &str, value: Buffer) -> Result<()> {
         let pool = self.get_client().await?;
 
-        let table = &self.table;
-        let key_field = &self.key_field;
-        let value_field = &self.value_field;
-        sqlx::query(&format!(
-            r#"INSERT INTO "{table}" ("{key_field}", "{value_field}")
-                VALUES ($1, $2)
-                ON CONFLICT ("{key_field}")
-                    DO UPDATE SET "{value_field}" = EXCLUDED."{value_field}""#,
-        ))
-        .bind(path)
-        .bind(value.to_vec())
-        .execute(pool)
-        .await
-        .map_err(parse_postgres_error)?;
+        sqlx::query(self.query_set())
+            .bind(path)
+            .bind(value.to_vec())
+            .execute(pool)
+            .await
+            .map_err(parse_postgres_error)?;
 
         Ok(())
     }
@@ -111,14 +179,11 @@ impl PostgresqlCore {
     pub async fn delete(&self, path: &str) -> Result<()> {
         let pool = self.get_client().await?;
 
-        sqlx::query(&format!(
-            r#"DELETE FROM "{}" WHERE "{}" = $1"#,
-            self.table, self.key_field
-        ))
-        .bind(path)
-        .execute(pool)
-        .await
-        .map_err(parse_postgres_error)?;
+        sqlx::query(self.query_delete())
+            .bind(path)
+            .execute(pool)
+            .await
+            .map_err(parse_postgres_error)?;
 
         Ok(())
     }
diff --git a/core/services/sqlite/src/backend.rs 
b/core/services/sqlite/src/backend.rs
index b32a9164e..94c84907e 100644
--- a/core/services/sqlite/src/backend.rs
+++ b/core/services/sqlite/src/backend.rs
@@ -138,13 +138,13 @@ impl Builder for SqliteBuilder {
 
         let root = normalize_root(self.config.root.as_deref().unwrap_or("/"));
 
-        Ok(SqliteBackend::new(SqliteCore {
-            pool: OnceCell::new(),
+        Ok(SqliteBackend::new(SqliteCore::new(
+            OnceCell::new(),
             config,
             table,
             key_field,
             value_field,
-        })
+        ))
         .with_normalized_root(root))
     }
 }
@@ -367,13 +367,13 @@ mod test {
     }
 
     async fn build_backend() -> SqliteBackend {
-        let core = SqliteCore {
-            pool: build_client().await,
-            config: Default::default(),
-            table: "test_table".to_string(),
-            key_field: "key".to_string(),
-            value_field: "value".to_string(),
-        };
+        let core = SqliteCore::new(
+            build_client().await,
+            Default::default(),
+            "test_table".to_string(),
+            "key".to_string(),
+            "value".to_string(),
+        );
 
         SqliteBackend::new(core)
     }
diff --git a/core/services/sqlite/src/core.rs b/core/services/sqlite/src/core.rs
index 478884595..a985af93e 100644
--- a/core/services/sqlite/src/core.rs
+++ b/core/services/sqlite/src/core.rs
@@ -16,6 +16,7 @@
 // under the License.
 
 use std::fmt::Debug;
+use std::sync::OnceLock;
 
 use mea::once::OnceCell;
 use sqlx::SqlitePool;
@@ -32,9 +33,39 @@ pub struct SqliteCore {
     pub table: String,
     pub key_field: String,
     pub value_field: String,
+
+    queries: SqliteQueries,
+}
+
+#[derive(Debug, Clone, Default)]
+struct SqliteQueries {
+    get: OnceLock<String>,
+    get_length: OnceLock<String>,
+    count_under: OnceLock<String>,
+    get_range_with_limit: OnceLock<String>,
+    get_range_without_limit: OnceLock<String>,
+    set: OnceLock<String>,
+    delete: OnceLock<String>,
 }
 
 impl SqliteCore {
+    pub(crate) fn new(
+        pool: OnceCell<SqlitePool>,
+        config: SqliteConnectOptions,
+        table: String,
+        key_field: String,
+        value_field: String,
+    ) -> Self {
+        Self {
+            pool,
+            config,
+            table,
+            key_field,
+            value_field,
+            queries: SqliteQueries::default(),
+        }
+    }
+
     pub async fn get_client(&self) -> Result<&SqlitePool> {
         self.pool
             .get_or_try_init(|| async {
@@ -45,6 +76,90 @@ impl SqliteCore {
             .await
     }
 
+    fn query_get(&self) -> &str {
+        self.queries
+            .get
+            .get_or_init(|| {
+                format!(
+                    r#"SELECT "{}" FROM "{}" WHERE "{}" = $1 LIMIT 1"#,
+                    self.value_field, self.table, self.key_field
+                )
+            })
+            .as_str()
+    }
+
+    fn query_get_length(&self) -> &str {
+        self.queries
+            .get_length
+            .get_or_init(|| {
+                format!(
+                    r#"SELECT LENGTH(CAST("{}" AS BLOB)) FROM "{}" WHERE "{}" 
= $1 LIMIT 1"#,
+                    self.value_field, self.table, self.key_field
+                )
+            })
+            .as_str()
+    }
+
+    fn query_count_under(&self) -> &str {
+        self.queries
+            .count_under
+            .get_or_init(|| {
+                format!(
+                    r#"SELECT COUNT(*) FROM "{}" WHERE "{}" LIKE $1 LIMIT 1"#,
+                    self.table, self.key_field
+                )
+            })
+            .as_str()
+    }
+
+    fn query_get_range_with_limit(&self) -> &str {
+        self.queries
+            .get_range_with_limit
+            .get_or_init(|| {
+                format!(
+                    r#"SELECT SUBSTR(CAST("{}" AS BLOB), $1, $2), 
LENGTH(CAST("{}" AS BLOB)) FROM "{}" WHERE "{}" = $3 LIMIT 1"#,
+                    self.value_field, self.value_field, self.table, 
self.key_field
+                )
+            })
+            .as_str()
+    }
+
+    fn query_get_range_without_limit(&self) -> &str {
+        self.queries
+            .get_range_without_limit
+            .get_or_init(|| {
+                format!(
+                    r#"SELECT SUBSTR(CAST("{}" AS BLOB), $1), LENGTH(CAST("{}" 
AS BLOB)) FROM "{}" WHERE "{}" = $2 LIMIT 1"#,
+                    self.value_field, self.value_field, self.table, 
self.key_field
+                )
+            })
+            .as_str()
+    }
+
+    fn query_set(&self) -> &str {
+        self.queries
+            .set
+            .get_or_init(|| {
+                format!(
+                    r#"INSERT OR REPLACE INTO "{}" ("{}", "{}") VALUES ($1, 
$2)"#,
+                    self.table, self.key_field, self.value_field,
+                )
+            })
+            .as_str()
+    }
+
+    fn query_delete(&self) -> &str {
+        self.queries
+            .delete
+            .get_or_init(|| {
+                format!(
+                    r#"DELETE FROM "{}" WHERE "{}" = $1"#,
+                    self.table, self.key_field
+                )
+            })
+            .as_str()
+    }
+
     pub async fn get(&self, path: &str) -> Result<Option<Buffer>> {
         let pool = self.get_client().await?;
 
@@ -57,14 +172,11 @@ impl SqliteCore {
         // - formatted, quoted identifiers for trusted table and field 
configuration
         // - bind parameters for values to avoid malformed SQL and SQL 
injection
         // to ensure correctness.
-        let value: Option<Vec<u8>> = sqlx::query_scalar(&format!(
-            r#"SELECT "{}" FROM "{}" WHERE "{}" = $1 LIMIT 1"#,
-            self.value_field, self.table, self.key_field
-        ))
-        .bind(path)
-        .fetch_optional(pool)
-        .await
-        .map_err(parse_sqlite_error)?;
+        let value: Option<Vec<u8>> = sqlx::query_scalar(self.query_get())
+            .bind(path)
+            .fetch_optional(pool)
+            .await
+            .map_err(parse_sqlite_error)?;
 
         Ok(value.map(Buffer::from))
     }
@@ -72,14 +184,11 @@ impl SqliteCore {
     pub async fn get_length(&self, path: &str) -> Result<Option<usize>> {
         let pool = self.get_client().await?;
 
-        let value: Option<i64> = sqlx::query_scalar(&format!(
-            r#"SELECT LENGTH(CAST("{}" AS BLOB)) FROM "{}" WHERE "{}" = $1 
LIMIT 1"#,
-            self.value_field, self.table, self.key_field
-        ))
-        .bind(path)
-        .fetch_optional(pool)
-        .await
-        .map_err(parse_sqlite_error)?;
+        let value: Option<i64> = sqlx::query_scalar(self.query_get_length())
+            .bind(path)
+            .fetch_optional(pool)
+            .await
+            .map_err(parse_sqlite_error)?;
 
         value
             .map(|v| {
@@ -94,14 +203,11 @@ impl SqliteCore {
     pub async fn count_under(&self, path: &str) -> Result<i64> {
         let pool = self.get_client().await?;
 
-        sqlx::query_scalar(&format!(
-            r#"SELECT COUNT(*) FROM "{}" WHERE "{}" LIKE $1 LIMIT 1"#,
-            self.table, self.key_field
-        ))
-        .bind(format!("{}%", path))
-        .fetch_one(pool)
-        .await
-        .map_err(parse_sqlite_error)
+        sqlx::query_scalar(self.query_count_under())
+            .bind(format!("{}%", path))
+            .fetch_one(pool)
+            .await
+            .map_err(parse_sqlite_error)
     }
 
     pub async fn get_range(
@@ -140,25 +246,19 @@ impl SqliteCore {
                     )
                     .set_source(err)
                 })?;
-                sqlx::query_as(&format!(
-                    r#"SELECT SUBSTR(CAST("{}" AS BLOB), $1, $2), 
LENGTH(CAST("{}" AS BLOB)) FROM "{}" WHERE "{}" = $3 LIMIT 1"#,
-                    self.value_field, self.value_field, self.table, 
self.key_field
-                ))
-                .bind(start)
-                .bind(limit)
-                .bind(path)
-                .fetch_optional(pool)
-                .await
+                sqlx::query_as(self.query_get_range_with_limit())
+                    .bind(start)
+                    .bind(limit)
+                    .bind(path)
+                    .fetch_optional(pool)
+                    .await
             }
             None => {
-                sqlx::query_as(&format!(
-                    r#"SELECT SUBSTR(CAST("{}" AS BLOB), $1), LENGTH(CAST("{}" 
AS BLOB)) FROM "{}" WHERE "{}" = $2 LIMIT 1"#,
-                    self.value_field, self.value_field, self.table, 
self.key_field
-                ))
-                .bind(start)
-                .bind(path)
-                .fetch_optional(pool)
-                .await
+                sqlx::query_as(self.query_get_range_without_limit())
+                    .bind(start)
+                    .bind(path)
+                    .fetch_optional(pool)
+                    .await
             }
         };
         let value: Option<(Vec<u8>, i64)> = value.map_err(parse_sqlite_error)?;
@@ -169,15 +269,12 @@ impl SqliteCore {
     pub async fn set(&self, path: &str, value: Buffer) -> Result<()> {
         let pool = self.get_client().await?;
 
-        sqlx::query(&format!(
-            r#"INSERT OR REPLACE INTO "{}" ("{}", "{}") VALUES ($1, $2)"#,
-            self.table, self.key_field, self.value_field,
-        ))
-        .bind(path)
-        .bind(value.to_vec())
-        .execute(pool)
-        .await
-        .map_err(parse_sqlite_error)?;
+        sqlx::query(self.query_set())
+            .bind(path)
+            .bind(value.to_vec())
+            .execute(pool)
+            .await
+            .map_err(parse_sqlite_error)?;
 
         Ok(())
     }
@@ -185,14 +282,11 @@ impl SqliteCore {
     pub async fn delete(&self, path: &str) -> Result<()> {
         let pool = self.get_client().await?;
 
-        sqlx::query(&format!(
-            r#"DELETE FROM "{}" WHERE "{}" = $1"#,
-            self.table, self.key_field
-        ))
-        .bind(path)
-        .execute(pool)
-        .await
-        .map_err(parse_sqlite_error)?;
+        sqlx::query(self.query_delete())
+            .bind(path)
+            .execute(pool)
+            .await
+            .map_err(parse_sqlite_error)?;
 
         Ok(())
     }

Reply via email to