From b4362726427aa7d6aaf29043e7ede0622c8ce7a6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=AE=89=E6=AD=A3=E8=B6=85?= Date: Fri, 29 May 2026 08:57:57 +0800 Subject: [PATCH] fix: tolerate stalled listing readers after quorum (#3110) Co-authored-by: loverustfs --- crates/e2e_test/src/lib.rs | 4 + .../src/mc_mirror_small_bucket_test.rs | 124 ++++++++++++++++++ .../ecstore/src/cache_value/metacache_set.rs | 37 +++++- 3 files changed, 159 insertions(+), 6 deletions(-) create mode 100644 crates/e2e_test/src/mc_mirror_small_bucket_test.rs diff --git a/crates/e2e_test/src/lib.rs b/crates/e2e_test/src/lib.rs index 96ce1bfa0..4fe216f03 100644 --- a/crates/e2e_test/src/lib.rs +++ b/crates/e2e_test/src/lib.rs @@ -63,6 +63,10 @@ mod archive_download_integrity_test; #[cfg(test)] mod list_objects_v2_pagination_test; +// Regression test for Issue #3107: mc mirror small-bucket listing must not time out. +#[cfg(test)] +mod mc_mirror_small_bucket_test; + // Policy variables tests #[cfg(test)] mod policy; diff --git a/crates/e2e_test/src/mc_mirror_small_bucket_test.rs b/crates/e2e_test/src/mc_mirror_small_bucket_test.rs new file mode 100644 index 000000000..6654a8663 --- /dev/null +++ b/crates/e2e_test/src/mc_mirror_small_bucket_test.rs @@ -0,0 +1,124 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use crate::common::{DEFAULT_ACCESS_KEY, DEFAULT_SECRET_KEY, RustFSTestEnvironment}; +use serial_test::serial; +use std::path::Path; +use std::process::Command; +use std::time::Duration; +use tokio::fs; +use tokio::time::timeout; +use tracing::info; + +const BUCKET: &str = "ddd"; +const OBJECT_COUNT: usize = 484; + +type TestResult = Result<(), Box>; + +async fn create_issue_3107_fixture(root: &Path) -> TestResult { + for idx in 0..OBJECT_COUNT { + let object_path = root.join("ddd").join("requirements").join(format!("file-{idx:04}.txt")); + fs::create_dir_all(object_path.parent().expect("object parent should exist")).await?; + fs::write(object_path, format!("issue-3107 payload {idx}\n").repeat(128)).await?; + } + + for dir in ["aaa", "ccc", "ccc-package"] { + let object_path = root.join(dir).join("marker.txt"); + fs::create_dir_all(object_path.parent().expect("object parent should exist")).await?; + fs::write(object_path, dir.as_bytes()).await?; + } + + Ok(()) +} + +fn mc_available() -> bool { + Command::new("mc") + .arg("--version") + .output() + .is_ok_and(|output| output.status.success()) +} + +fn run_mc(args: &[&str]) -> TestResult { + let output = Command::new("mc").args(args).output()?; + if !output.status.success() { + return Err(format!( + "mc {:?} failed: status={:?}, stdout={}, stderr={}", + args, + output.status.code(), + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ) + .into()); + } + Ok(()) +} + +fn count_files(root: &Path) -> usize { + walkdir::WalkDir::new(root) + .into_iter() + .filter_map(Result::ok) + .filter(|entry| entry.file_type().is_file()) + .count() +} + +#[tokio::test] +#[serial] +async fn test_mc_mirror_small_bucket_completes_without_list_timeout() -> TestResult { + crate::common::init_logging(); + info!("Starting issue #3107 mc mirror regression test"); + if !mc_available() { + info!("Skipping issue #3107 mc mirror regression test because mc is not installed"); + return Ok(()); + } + + let mut env = RustFSTestEnvironment::new().await?; + env.start_rustfs_server(vec![]).await?; + env.create_test_bucket(BUCKET).await?; + + let alias = format!("rustfs-issue-3107-{}", uuid::Uuid::new_v4()); + run_mc(&["alias", "set", &alias, &env.url, DEFAULT_ACCESS_KEY, DEFAULT_SECRET_KEY])?; + + let fixture = Path::new(&env.temp_dir).join("issue-3107-fixture"); + let backup = Path::new(&env.temp_dir).join("issue-3107-backup"); + fs::create_dir_all(&fixture).await?; + fs::create_dir_all(&backup).await?; + create_issue_3107_fixture(&fixture).await?; + + let bucket_target = format!("{alias}/{BUCKET}"); + let source_path = fixture.to_string_lossy().to_string(); + run_mc(&["mirror", "--overwrite", &source_path, &bucket_target])?; + + let backup_path = backup.join("ddd-backup"); + let backup_path_string = backup_path.to_string_lossy().to_string(); + timeout( + Duration::from_secs(20), + tokio::task::spawn_blocking({ + let bucket_target = bucket_target.clone(); + let backup_path_string = backup_path_string.clone(); + move || run_mc(&["mirror", "--overwrite", &bucket_target, &backup_path_string]) + }), + ) + .await + .map_err(|_| "mc mirror timed out while listing/comparing a small bucket")???; + + let mirrored_count = count_files(&backup_path); + assert_eq!( + mirrored_count, + OBJECT_COUNT + 3, + "mc mirror should copy every object from the issue #3107 small bucket fixture" + ); + + run_mc(&["alias", "remove", &alias])?; + Ok(()) +} diff --git a/crates/ecstore/src/cache_value/metacache_set.rs b/crates/ecstore/src/cache_value/metacache_set.rs index 0718c470a..8910d896d 100644 --- a/crates/ecstore/src/cache_value/metacache_set.rs +++ b/crates/ecstore/src/cache_value/metacache_set.rs @@ -442,10 +442,13 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d { finished_fn(&errs).await; } - if errs.iter().flatten().any(|err| *err == DiskError::Timeout) { - return Err(DiskError::Timeout); - } - + // All remaining readers are either at EOF or failed. Earlier logic + // returned Timeout here for even a single stalled drive, despite + // `has_err` being within the tolerated drive-failure budget. That + // makes small distributed listings fail once healthy quorum readers + // reach EOF but one remote walk stream is slow/stalled. Only the + // `has_err > opts.disks.len() - opts.min_disks` branch above should + // turn tolerated reader failures into request failures. // error!("list_path_raw: at_eof + has_err == readers.len() break {:?}", &errs); break; } @@ -486,6 +489,12 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d return Err(err); } + // The merge consumer can finish successfully before every producer finishes + // (for example after reaching EOF quorum while a tolerated drive is stalled, + // or after the requested listing limit is satisfied). Cancel remaining walk + // jobs before joining them so list calls do not wait for slow remote streams. + cancel_rx.cancel(); + let results = join_all(jobs).await; let mut job_errs = Vec::new(); for result in results { @@ -545,18 +554,34 @@ mod tests { CancellationToken::new(), ListPathRawOptions { disks: vec![None, None], - min_disks: 1, + min_disks: 2, test_reader_behaviors: vec![TestReaderBehavior::Stall, TestReaderBehavior::Eof], peek_timeout: Some(Duration::from_millis(20)), ..Default::default() }, ) .await - .expect_err("stalled reader should make listing fail explicitly"); + .expect_err("stalled reader should fail when read quorum cannot be met"); assert_eq!(err, DiskError::Timeout); } + #[tokio::test] + async fn list_path_raw_tolerates_stalled_reader_after_quorum_eof() { + list_path_raw( + CancellationToken::new(), + ListPathRawOptions { + disks: vec![None, None, None], + min_disks: 2, + test_reader_behaviors: vec![TestReaderBehavior::Eof, TestReaderBehavior::Eof, TestReaderBehavior::Stall], + peek_timeout: Some(Duration::from_millis(20)), + ..Default::default() + }, + ) + .await + .expect("listing should complete when healthy quorum reached EOF and only a tolerated drive stalled"); + } + #[tokio::test] async fn list_path_raw_returns_timeout_when_producer_fails_after_partial_entry() { let seen = Arc::new(Mutex::new(Vec::new()));