This is an automated email from the ASF dual-hosted git repository. erickguan pushed a commit to branch fix-sqlite-stat in repository https://gitbox.apache.org/repos/asf/opendal.git
commit f465ebd70419505940d0b397df62a4231658f093 Author: Erick Guan <[email protected]> AuthorDate: Sat Jul 4 17:18:01 2026 +0800 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,
