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();