This is an automated email from the ASF dual-hosted git repository.
erickguan pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/opendal.git
The following commit(s) were added to refs/heads/main by this push:
new 140b2de0a fix(services/sqlite): stat getting the entire object (#7853)
140b2de0a is described below
commit 140b2de0a6307ca4fdccde8650c110148c6d9ebf
Author: Erick Guan <[email protected]>
AuthorDate: Sat Jul 4 21:30:52 2026 +0800
fix(services/sqlite): stat getting the entire object (#7853)
Fix sqlite stat getting the entire object
---
core/services/sqlite/src/backend.rs | 155 +++++++++++++++++++++++++++---------
core/services/sqlite/src/core.rs | 6 +-
2 files changed, 120 insertions(+), 41 deletions(-)
diff --git a/core/services/sqlite/src/backend.rs
b/core/services/sqlite/src/backend.rs
index 63feabbac..5a96dae0e 100644
--- a/core/services/sqlite/src/backend.rs
+++ b/core/services/sqlite/src/backend.rs
@@ -228,10 +228,10 @@ impl Service for SqliteBackend {
if p == build_abs_path(&self.root, "") {
Ok(RpStat::new(Metadata::new(EntryMode::DIR)))
} else {
- let bs = self.core.get(&p).await?;
- match bs {
- Some(bs) => Ok(RpStat::new(
-
Metadata::new(EntryMode::from_path(&p)).with_content_length(bs.len() as u64),
+ let length = self.core.get_length(&p).await?;
+ match length {
+ Some(length) => Ok(RpStat::new(
+
Metadata::new(EntryMode::from_path(&p)).with_content_length(length as u64),
)),
None => {
// Check if this might be a directory by looking for keys
with this prefix
@@ -373,68 +373,147 @@ mod test {
OnceCell::from_value(pool)
}
- #[tokio::test]
- async fn test_sqlite_accessor_creation() {
+ async fn build_backend() -> SqliteBackend {
let core = SqliteCore {
pool: build_client().await,
config: Default::default(),
- table: "test".to_string(),
+ table: "test_table".to_string(),
key_field: "key".to_string(),
value_field: "value".to_string(),
};
- let accessor = SqliteBackend::new(core);
+ SqliteBackend::new(core)
+ }
+
+ #[tokio::test]
+ async fn test_sqlite_backend_creation() {
+ let backend = build_backend().await;
// Verify basic properties
- assert_eq!(accessor.root, "/");
- assert_eq!(accessor.info.scheme(), SQLITE_SCHEME);
- assert!(accessor.capability().read);
- assert!(accessor.capability().write);
- assert!(accessor.capability().delete);
- assert!(accessor.capability().stat);
+ assert_eq!(backend.root, "/");
+ assert_eq!(backend.info.scheme(), SQLITE_SCHEME);
+ assert!(backend.capability().read);
+ assert!(backend.capability().write);
+ assert!(backend.capability().delete);
+ assert!(backend.capability().stat);
}
#[tokio::test]
- async fn test_sqlite_accessor_with_root() {
- let core = SqliteCore {
- pool: build_client().await,
- config: Default::default(),
- table: "test".to_string(),
- key_field: "key".to_string(),
- value_field: "value".to_string(),
- };
-
- let accessor =
SqliteBackend::new(core).with_normalized_root("/test/".to_string());
+ async fn test_sqlite_backend_with_root() {
+ let backend = build_backend()
+ .await
+ .with_normalized_root("/test/".to_string());
- assert_eq!(accessor.root, "/test/");
- assert_eq!(accessor.info.root(), Arc::from("/test/"));
+ assert_eq!(backend.root, "/test/");
+ assert_eq!(backend.info.root(), Arc::from("/test/"));
}
#[tokio::test]
async fn test_sqlite_read_range_from_offset_reads_to_eof() {
- let core = SqliteCore {
- pool: build_client().await,
- config: Default::default(),
- table: "test".to_string(),
- key_field: "key".to_string(),
- value_field: "value".to_string(),
- };
- let pool = core.get_client().await.unwrap();
- sqlx::query("CREATE TABLE test (key TEXT PRIMARY KEY, value BLOB)")
+ let backend = build_backend().await;
+
+ let pool = backend.core.get_client().await.unwrap();
+ sqlx::query("CREATE TABLE test_table (key TEXT PRIMARY KEY, value
BLOB)")
.execute(pool)
.await
.unwrap();
- let accessor = SqliteBackend::new(core);
let ctx = OperationContext::new();
- let mut writer = accessor.write(&ctx, "hello",
OpWrite::default()).unwrap();
+ let mut writer = backend.write(&ctx, "hello",
OpWrite::default()).unwrap();
writer.write(Buffer::from("hello world")).await.unwrap();
writer.close().await.unwrap();
- let reader = accessor.read(&ctx, "hello", OpRead::default()).unwrap();
+ let reader = backend.read(&ctx, "hello", OpRead::default()).unwrap();
let (_, mut stream) =
reader.open(BytesRange::from(6_u64..)).await.unwrap();
let buffer = stream.read_all().await.unwrap();
assert_eq!(buffer.to_vec(), b"world");
}
+
+ #[tokio::test]
+ async fn test_sqlite_stat_uses_value_length() {
+ let backend = build_backend().await;
+
+ let pool = backend.core.get_client().await.unwrap();
+ sqlx::query("CREATE TABLE test_table (key TEXT PRIMARY KEY, value
BLOB)")
+ .execute(pool)
+ .await
+ .unwrap();
+
+ let ctx = OperationContext::new();
+ let mut writer = backend.write(&ctx, "key_id",
OpWrite::default()).unwrap();
+ writer.write(Buffer::from("hello world")).await.unwrap();
+ writer.close().await.unwrap();
+
+ let rp = backend
+ .stat(&ctx, "key_id", OpStat::default())
+ .await
+ .unwrap();
+
+ assert_eq!(rp.into_metadata().content_length(), 11);
+ }
+
+ #[tokio::test]
+ async fn test_sqlite_stat_returns_byte_length_for_text_value() {
+ let backend = build_backend().await;
+
+ let pool = backend.core.get_client().await.unwrap();
+ sqlx::query("CREATE TABLE test_table (key TEXT PRIMARY KEY, value
BLOB)")
+ .execute(pool)
+ .await
+ .unwrap();
+ sqlx::query("INSERT INTO test_table (key, value) VALUES ($1, $2)")
+ .bind("key_id")
+ .bind("你好")
+ .execute(pool)
+ .await
+ .unwrap();
+
+ let ctx = OperationContext::new();
+
+ let rp = backend
+ .stat(&ctx, "key_id", OpStat::default())
+ .await
+ .unwrap();
+ assert_eq!(rp.into_metadata().content_length(), 6);
+
+ let reader = backend.read(&ctx, "key_id", OpRead::default()).unwrap();
+ let (rp, mut stream) =
reader.open(BytesRange::from(0_u64..3)).await.unwrap();
+ let buffer = stream.read_all().await.unwrap();
+
+ assert_eq!(rp.into_metadata().unwrap().content_length(), 6);
+ assert_eq!(buffer.to_vec(), "你".as_bytes());
+ }
+
+ #[tokio::test]
+ async fn test_sqlite_stat_returns_byte_length_for_text_column() {
+ let backend = build_backend().await;
+ let pool = backend.core.get_client().await.unwrap();
+
+ sqlx::query("CREATE TABLE test_table (key TEXT PRIMARY KEY, value
TEXT)")
+ .execute(pool)
+ .await
+ .unwrap();
+ sqlx::query("INSERT INTO test_table (key, value) VALUES ($1, $2)")
+ .bind("key_id")
+ .bind("你好")
+ .execute(pool)
+ .await
+ .unwrap();
+
+ let ctx = OperationContext::new();
+
+ let rp = backend
+ .stat(&ctx, "key_id", OpStat::default())
+ .await
+ .unwrap();
+ assert_eq!(rp.into_metadata().content_length(), 6);
+
+ let reader = backend.read(&ctx, "key_id", OpRead::default()).unwrap();
+ let (rp, mut stream) =
reader.open(BytesRange::from(0_u64..3)).await.unwrap();
+ let buffer = stream.read_all().await.unwrap();
+
+ assert_eq!(rp.into_metadata().unwrap().content_length(), 6);
+ assert_eq!(buffer.to_vec(), "你".as_bytes());
+ }
}
diff --git a/core/services/sqlite/src/core.rs b/core/services/sqlite/src/core.rs
index 9cc82440d..16ee5f68a 100644
--- a/core/services/sqlite/src/core.rs
+++ b/core/services/sqlite/src/core.rs
@@ -64,7 +64,7 @@ impl SqliteCore {
let pool = self.get_client().await?;
let value: Option<i64> = sqlx::query_scalar(&format!(
- "SELECT LENGTH(`{}`) FROM `{}` WHERE `{}` = $1 LIMIT 1",
+ "SELECT LENGTH(CAST(`{}` AS BLOB)) FROM `{}` WHERE `{}` = $1 LIMIT
1",
self.value_field, self.table, self.key_field
))
.bind(path)
@@ -91,7 +91,7 @@ impl SqliteCore {
let pool = self.get_client().await?;
let query = match limit {
Some(limit) => format!(
- "SELECT SUBSTR(`{}`, {}, {}), LENGTH(`{}`) FROM `{}` WHERE
`{}` = $1 LIMIT 1",
+ "SELECT SUBSTR(CAST(`{}` AS BLOB), {}, {}), LENGTH(CAST(`{}`
AS BLOB)) FROM `{}` WHERE `{}` = $1 LIMIT 1",
self.value_field,
start + 1,
limit,
@@ -100,7 +100,7 @@ impl SqliteCore {
self.key_field
),
None => format!(
- "SELECT SUBSTR(`{}`, {}), LENGTH(`{}`) FROM `{}` WHERE `{}` =
$1 LIMIT 1",
+ "SELECT SUBSTR(CAST(`{}` AS BLOB), {}), LENGTH(CAST(`{}` AS
BLOB)) FROM `{}` WHERE `{}` = $1 LIMIT 1",
self.value_field,
start + 1,
self.value_field,