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 5160785c7 fix(layers/foyer): preserve read args while filling 
full-object cache (#7718)
5160785c7 is described below

commit 5160785c740fdb7ea0961bdee3b076edd16340b9
Author: tonghuaroot (童话) <[email protected]>
AuthorDate: Sun Jun 28 23:55:37 2026 +0800

    fix(layers/foyer): preserve read args while filling full-object cache 
(#7718)
    
    The full-object cache fill called `stat` and `read` with default args,
    so versioned and conditional reads stored latest content under a
    versioned cache key. Derive the stat args from the reader's `OpRead`
    and pass the original read args, matching the fallback paths.
    
    Closes #7686
    
    Co-authored-by: tonghuaroot <[email protected]>
---
 core/layers/foyer/src/full.rs | 28 ++++++++++++++++++++++++++--
 core/layers/foyer/src/lib.rs  | 34 ++++++++++++++++++++++++++++++++--
 2 files changed, 58 insertions(+), 4 deletions(-)

diff --git a/core/layers/foyer/src/full.rs b/core/layers/foyer/src/full.rs
index bdcb39a8a..0325d737c 100644
--- a/core/layers/foyer/src/full.rs
+++ b/core/layers/foyer/src/full.rs
@@ -81,6 +81,28 @@ impl FullReader {
             .contains(key)
     }
 
+    /// Derive the `OpStat` used for cache fill from the reader's `OpRead`, so 
the
+    /// stat observes the same conditions (version, if-match, ...) as the 
request.
+    fn stat_args(&self) -> OpStat {
+        let mut op = OpStat::new();
+        if let Some(v) = self.args.version() {
+            op = op.with_version(v);
+        }
+        if let Some(v) = self.args.if_match() {
+            op = op.with_if_match(v);
+        }
+        if let Some(v) = self.args.if_none_match() {
+            op = op.with_if_none_match(v);
+        }
+        if let Some(v) = self.args.if_modified_since() {
+            op = op.with_if_modified_since(v);
+        }
+        if let Some(v) = self.args.if_unmodified_since() {
+            op = op.with_if_unmodified_since(v);
+        }
+        op
+    }
+
     async fn fallback_open(
         &self,
         range: BytesRange,
@@ -123,11 +145,13 @@ impl FullReader {
                 let inner = self.inner.clone();
                 let size_limit = self.size_limit.clone();
                 let path_clone = path_str.clone();
+                let stat_args = self.stat_args();
+                let read_args = self.args.clone();
                 async move {
                     // read the metadata first, if it's too large, do not cache
                     let metadata = inner
                         .srv
-                        .stat(&inner.ctx, &path_clone, OpStat::default())
+                        .stat(&inner.ctx, &path_clone, stat_args)
                         .await
                         .map_err(FetchError::from_error)?
                         .into_metadata();
@@ -140,7 +164,7 @@ impl FullReader {
                     // fetch the ENTIRE object from remote.
                     let reader = inner
                         .srv
-                        .read(&inner.ctx, &path_clone, OpRead::default())
+                        .read(&inner.ctx, &path_clone, read_args)
                         .map_err(FetchError::from_error)?;
                     let (_, mut stream) = reader
                         .open(BytesRange::new(0, None))
diff --git a/core/layers/foyer/src/lib.rs b/core/layers/foyer/src/lib.rs
index b73a2e3f9..9dbaf81ff 100644
--- a/core/layers/foyer/src/lib.rs
+++ b/core/layers/foyer/src/lib.rs
@@ -382,6 +382,8 @@ mod tests {
         stat_calls: AtomicUsize,
         open_calls: AtomicUsize,
         read_calls: AtomicUsize,
+        last_stat_args: Mutex<Option<OpStat>>,
+        last_read_args: Mutex<Option<OpRead>>,
     }
 
     impl MockReadState {
@@ -391,6 +393,8 @@ mod tests {
                 stat_calls: AtomicUsize::new(0),
                 open_calls: AtomicUsize::new(0),
                 read_calls: AtomicUsize::new(0),
+                last_stat_args: Mutex::new(None),
+                last_read_args: Mutex::new(None),
             }
         }
 
@@ -479,12 +483,14 @@ mod tests {
             ))
         }
 
-        async fn stat(&self, _: &OperationContext, _: &str, _: OpStat) -> 
Result<RpStat> {
+        async fn stat(&self, _: &OperationContext, _: &str, args: OpStat) -> 
Result<RpStat> {
             self.state.stat_calls.fetch_add(1, Ordering::Relaxed);
+            *self.state.last_stat_args.lock().unwrap() = Some(args);
             Ok(RpStat::new(self.state.metadata()))
         }
 
-        fn read(&self, _ctx: &OperationContext, _: &str, _: OpRead) -> 
Result<Self::Reader> {
+        fn read(&self, _ctx: &OperationContext, _: &str, args: OpRead) -> 
Result<Self::Reader> {
+            *self.state.last_read_args.lock().unwrap() = Some(args);
             Ok(MockReadReader {
                 state: self.state.clone(),
             })
@@ -599,6 +605,30 @@ mod tests {
         assert_eq!(state.read_calls.load(Ordering::Relaxed), 0);
     }
 
+    #[tokio::test]
+    async fn test_cache_fill_preserves_read_args() {
+        let cache = memory_cache().await;
+        let source = Arc::new(MockReadService::new("0123456789"));
+        let state = source.state.clone();
+        let service = FoyerLayer::new(cache)
+            .with_size_limit(0..100)
+            .apply_service(source);
+        let ctx = service_context(&service);
+
+        let args = 
OpRead::default().with_version("v1").with_if_match("etag-1");
+        let reader = service.read(&ctx, "test", args).unwrap();
+        let (_, mut stream) = reader.open(BytesRange::new(0, 
None)).await.unwrap();
+        stream.read_all().await.unwrap();
+
+        let stat_args = state.last_stat_args.lock().unwrap().clone().unwrap();
+        assert_eq!(stat_args.version(), Some("v1"));
+        assert_eq!(stat_args.if_match(), Some("etag-1"));
+
+        let read_args = state.last_read_args.lock().unwrap().clone().unwrap();
+        assert_eq!(read_args.version(), Some("v1"));
+        assert_eq!(read_args.if_match(), Some("etag-1"));
+    }
+
     #[tokio::test]
     async fn test() {
         let dir = tempfile::tempdir().unwrap();

Reply via email to