From 26341742e05624d3338bf2c38c01002921ebcab2 Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Wed, 13 May 2026 19:08:05 +0800 Subject: [PATCH] fix(ecstore): fail listing on walk_dir producer errors (#2937) --- .../ecstore/src/cache_value/metacache_set.rs | 149 +++++++++++------- crates/ecstore/src/disk/disk_store.rs | 2 +- crates/ecstore/src/rpc/remote_disk.rs | 2 +- 3 files changed, 96 insertions(+), 57 deletions(-) diff --git a/crates/ecstore/src/cache_value/metacache_set.rs b/crates/ecstore/src/cache_value/metacache_set.rs index a78d7d8ee..d69a8b596 100644 --- a/crates/ecstore/src/cache_value/metacache_set.rs +++ b/crates/ecstore/src/cache_value/metacache_set.rs @@ -18,7 +18,12 @@ use crate::disk::{self, DiskAPI, DiskStore, WalkDirOptions}; use futures::future::join_all; use metrics::counter; use rustfs_filemeta::{MetaCacheEntries, MetaCacheEntry, MetacacheReader, is_io_eof}; -use std::{future::Future, pin::Pin, time::Duration}; +use std::{ + future::Future, + pin::Pin, + sync::{Arc, Mutex}, + time::Duration, +}; use tokio::io::AsyncRead; use tokio::spawn; use tokio::time::timeout; @@ -46,10 +51,12 @@ async fn peek_with_timeout(reader: &mut MetacacheReader } #[cfg(test)] -#[derive(Clone, Copy)] +#[derive(Clone)] pub(crate) enum TestReaderBehavior { Eof, Stall, + ProducerError(DiskError), + PartialThenTimeout(Vec), } #[derive(Default)] @@ -107,21 +114,21 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d let mut readers = Vec::with_capacity(opts.disks.len()); let fds = opts.fallback_disks.iter().flatten().cloned().collect::>(); let max_disk_failures = opts.disks.len().saturating_sub(opts.min_disks); + let producer_errs = Arc::new(Mutex::new(vec![None; opts.disks.len()])); let cancel_rx = CancellationToken::new(); for (disk_idx, disk) in opts.disks.iter().enumerate() { - #[cfg(not(test))] - let _ = disk_idx; let opdisk = disk.clone(); let opts_clone = opts.clone(); let mut fds_clone = fds.clone(); let cancel_rx_clone = cancel_rx.clone(); + let producer_errs_clone = producer_errs.clone(); let (rd, wr) = tokio::io::duplex(64); readers.push(MetacacheReader::new(rd)); jobs.push(spawn(async move { #[cfg(test)] - if let Some(behavior) = opts_clone.test_reader_behaviors.get(disk_idx).copied() { + if let Some(behavior) = opts_clone.test_reader_behaviors.get(disk_idx).cloned() { match behavior { TestReaderBehavior::Eof => return Ok(()), TestReaderBehavior::Stall => { @@ -129,6 +136,19 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d cancel_rx_clone.cancelled().await; return Ok(()); } + TestReaderBehavior::ProducerError(err) => { + producer_errs_clone.lock().expect("producer error mutex poisoned")[disk_idx] = Some(err.clone()); + return Err(err); + } + TestReaderBehavior::PartialThenTimeout(entries) => { + let mut wr = wr; + let mut out = rustfs_filemeta::MetacacheWriter::new(&mut wr); + let err = DiskError::Timeout; + producer_errs_clone.lock().expect("producer error mutex poisoned")[disk_idx] = Some(err.clone()); + let _ = out.write(&entries).await; + drop(out); + return Err(err); + } } } @@ -166,18 +186,20 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d } while need_fallback { - let disk_op = { - if fds_clone.is_empty() { - None - } else { - let disk = fds_clone.remove(0); - if disk.is_online().await { Some(disk.clone()) } else { None } + let mut disk_op = None; + while !fds_clone.is_empty() { + let disk = fds_clone.remove(0); + if disk.is_online().await { + disk_op = Some(disk); + break; } - }; + } let Some(disk) = disk_op else { warn!("list_path_raw: fallback disk is none"); - break; + let err = last_err.unwrap_or(DiskError::DiskNotFound); + producer_errs_clone.lock().expect("producer error mutex poisoned")[disk_idx] = Some(err.clone()); + return Err(err); }; match disk @@ -204,7 +226,6 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d Err(err) => { error!("walk dir2 err {:?}", &err); last_err = Some(err); - break; } } } @@ -260,6 +281,11 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d // info!("read entry disk: {}, name: {}", i, entry.name); entry } else { + if let Some(err) = producer_errs.lock().expect("producer error mutex poisoned")[i].clone() { + has_err += 1; + errs[i] = Some(err); + continue; + } // eof at_eof += 1; // warn!("list_path_raw: peek eof, disk: {}", i); @@ -267,6 +293,12 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d } } PeekOutcome::Error(err) => { + if let Some(err) = producer_errs.lock().expect("producer error mutex poisoned")[i].clone() { + has_err += 1; + errs[i] = Some(err); + continue; + } + if err == rustfs_filemeta::Error::Unexpected { at_eof += 1; // warn!("list_path_raw: peek err eof, disk: {}", i); @@ -376,6 +408,15 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d if let Some(finished_fn) = opts.finished.as_ref() { finished_fn(&errs).await; } + if errs.iter().flatten().any(|err| *err == DiskError::Timeout) { + return Err(DiskError::Timeout); + } + let mut err_iter = errs.iter().flatten(); + if let Some(err) = err_iter.next() + && err_iter.next().is_none() + { + return Err(err.clone()); + } let mut combined_err = Vec::new(); errs.iter().zip(opts.disks.iter()).for_each(|(err, disk)| match (err, disk) { (Some(err), Some(disk)) => { @@ -451,7 +492,7 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d match result { Ok(Ok(())) => {} Ok(Err(err)) => { - error!("list_path_raw err {:?}", err); + error!("list_path_raw producer err {:?}", err); job_errs.push(err); } Err(err) => { @@ -472,8 +513,6 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d #[cfg(test)] mod tests { use super::*; - use crate::disk::endpoint::Endpoint; - use crate::disk::{DiskOption, new_disk}; use rustfs_filemeta::MetacacheWriter; #[tokio::test] @@ -503,6 +542,38 @@ mod tests { assert_eq!(err, DiskError::Timeout); } + #[tokio::test] + async fn list_path_raw_returns_timeout_when_producer_fails_after_partial_entry() { + let seen = Arc::new(Mutex::new(Vec::new())); + let seen_clone = seen.clone(); + + let err = list_path_raw( + CancellationToken::new(), + ListPathRawOptions { + disks: vec![None], + min_disks: 1, + test_reader_behaviors: vec![TestReaderBehavior::PartialThenTimeout(vec![MetaCacheEntry { + name: "bucket/object".to_string(), + metadata: vec![1, 2, 3], + cached: None, + reusable: false, + }])], + agreed: Some(Box::new(move |entry: MetaCacheEntry| { + let seen = seen_clone.clone(); + Box::pin(async move { + seen.lock().expect("seen mutex poisoned").push(entry.name); + }) + })), + ..Default::default() + }, + ) + .await + .expect_err("producer timeout after partial output must fail the listing"); + + assert_eq!(err, DiskError::Timeout); + assert_eq!(seen.lock().expect("seen mutex poisoned").as_slice(), &["bucket/object".to_string()]); + } + #[tokio::test] async fn peek_with_timeout_times_out_on_silent_reader() { let (_writer, reader) = tokio::io::duplex(64); @@ -536,52 +607,20 @@ mod tests { } } - #[cfg(unix)] #[tokio::test] - async fn list_path_raw_propagates_parent_scan_access_denied() { - use std::fs::Permissions; - use std::os::unix::fs::PermissionsExt; - use tempfile::tempdir; - use tokio::fs; - - let dir = tempdir().expect("tempdir should be created"); - let bucket = "test-bucket"; - let parent = dir.path().join(bucket).join("shiplog"); - - fs::create_dir_all(parent.join("nano/a.txt")) - .await - .expect("test object directory should be created"); - fs::write(parent.join("nano/a.txt/xl.meta"), b"meta") - .await - .expect("test metadata should be written"); - - std::fs::set_permissions(&parent, Permissions::from_mode(0o111)).expect("parent permissions should be changed"); - if fs::read_dir(&parent).await.is_ok() { - std::fs::set_permissions(&parent, Permissions::from_mode(0o755)).expect("parent permissions should be restored"); - return; - } - - let endpoint = Endpoint::try_from(dir.path().to_str().expect("temp path should be valid utf-8")) - .expect("endpoint should be created"); - let disk = new_disk(&endpoint, &DiskOption::default()) - .await - .expect("local disk should be created"); - - let result = list_path_raw( + async fn list_path_raw_propagates_producer_access_denied() { + let err = list_path_raw( CancellationToken::new(), ListPathRawOptions { - disks: vec![Some(disk)], - bucket: bucket.to_string(), - path: "shiplog/".to_string(), + disks: vec![None], min_disks: 1, + test_reader_behaviors: vec![TestReaderBehavior::ProducerError(DiskError::FileAccessDenied)], ..Default::default() }, ) - .await; + .await + .expect_err("producer access failure must not be treated as an empty listing"); - std::fs::set_permissions(&parent, Permissions::from_mode(0o755)).expect("parent permissions should be restored"); - - let err = result.expect_err("parent directory access failure must not be treated as an empty listing"); assert_eq!(err, DiskError::FileAccessDenied); } } diff --git a/crates/ecstore/src/disk/disk_store.rs b/crates/ecstore/src/disk/disk_store.rs index 6a7b6afe9..73a4da288 100644 --- a/crates/ecstore/src/disk/disk_store.rs +++ b/crates/ecstore/src/disk/disk_store.rs @@ -863,7 +863,7 @@ impl LocalDiskWrapper { timeout_ms = timeout_duration.as_millis(), "Local disk operation timed out" ); - Err(DiskError::other(format!("disk operation timeout after {timeout_duration:?}"))) + Err(DiskError::Timeout) } } } diff --git a/crates/ecstore/src/rpc/remote_disk.rs b/crates/ecstore/src/rpc/remote_disk.rs index c67a2c392..6859999ee 100644 --- a/crates/ecstore/src/rpc/remote_disk.rs +++ b/crates/ecstore/src/rpc/remote_disk.rs @@ -385,7 +385,7 @@ impl RemoteDisk { timeout_ms = timeout_duration.as_millis(), "Remote disk operation timed out" ); - Err(Error::other(format!("Remote disk operation timeout after {timeout_duration:?}"))) + Err(DiskError::Timeout) } } }