From 5ca99a8e32e24ea88e3d75b7f0942bcda29e1451 Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Mon, 11 May 2026 23:19:46 +0800 Subject: [PATCH] fix(ecstore): fail listing on stalled reader (#2920) Co-authored-by: houseme --- .../ecstore/src/cache_value/metacache_set.rs | 58 ++++++++++++++++++- 1 file changed, 56 insertions(+), 2 deletions(-) diff --git a/crates/ecstore/src/cache_value/metacache_set.rs b/crates/ecstore/src/cache_value/metacache_set.rs index c65a67275..251733c3b 100644 --- a/crates/ecstore/src/cache_value/metacache_set.rs +++ b/crates/ecstore/src/cache_value/metacache_set.rs @@ -45,6 +45,13 @@ async fn peek_with_timeout(reader: &mut MetacacheReader } } +#[cfg(test)] +#[derive(Clone, Copy)] +pub(crate) enum TestReaderBehavior { + Eof, + Stall, +} + #[derive(Default)] pub struct ListPathRawOptions { pub disks: Vec>, @@ -60,6 +67,10 @@ pub struct ListPathRawOptions { pub agreed: Option, pub partial: Option, pub finished: Option, + #[cfg(test)] + pub(crate) test_reader_behaviors: Vec, + #[cfg(test)] + pub(crate) peek_timeout: Option, // pub agreed: Option>, // pub partial: Option]) + Send + Sync>>, // pub finished: Option]) + Send + Sync>>, @@ -78,6 +89,10 @@ impl Clone for ListPathRawOptions { min_disks: self.min_disks, report_not_found: self.report_not_found, per_disk_limit: self.per_disk_limit, + #[cfg(test)] + test_reader_behaviors: self.test_reader_behaviors.clone(), + #[cfg(test)] + peek_timeout: self.peek_timeout, ..Default::default() } } @@ -94,14 +109,29 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d let cancel_rx = CancellationToken::new(); - for disk in opts.disks.iter() { + 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 (rd, mut wr) = tokio::io::duplex(64); + 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() { + match behavior { + TestReaderBehavior::Eof => return Ok(()), + TestReaderBehavior::Stall => { + let _held_writer = wr; + cancel_rx_clone.cancelled().await; + return Ok(()); + } + } + } + + let mut wr = wr; let wakl_opts = WalkDirOptions { bucket: opts_clone.bucket.clone(), base_dir: opts_clone.path.clone(), @@ -179,6 +209,9 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d } let revjob = spawn(async move { + #[cfg(test)] + let peek_timeout = opts.peek_timeout.unwrap_or_else(get_drive_walkdir_stall_timeout); + #[cfg(not(test))] let peek_timeout = get_drive_walkdir_stall_timeout(); let mut errs: Vec> = Vec::with_capacity(readers.len()); for _ in 0..readers.len() { @@ -358,6 +391,9 @@ 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); + } // error!("list_path_raw: at_eof + has_err == readers.len() break {:?}", &errs); break; @@ -424,6 +460,24 @@ mod tests { assert_eq!(err, DiskError::ErasureReadQuorum); } + #[tokio::test] + async fn list_path_raw_returns_timeout_when_reader_stalls_before_completion() { + let err = list_path_raw( + CancellationToken::new(), + ListPathRawOptions { + disks: vec![None, None], + min_disks: 1, + 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"); + + assert_eq!(err, DiskError::Timeout); + } + #[tokio::test] async fn peek_with_timeout_times_out_on_silent_reader() { let (_writer, reader) = tokio::io::duplex(64);