From 4c9431704cafee6dfdc5f31ed66dcc7f33169a74 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Sun, 12 Jul 2026 15:21:21 +0800 Subject: [PATCH] fix(ecstore): cancel orphaned listing walks (#4773) --- .../ecstore/src/cache_value/metacache_set.rs | 373 ++++++++++++++++-- crates/ecstore/src/store/list_objects.rs | 18 +- 2 files changed, 352 insertions(+), 39 deletions(-) diff --git a/crates/ecstore/src/cache_value/metacache_set.rs b/crates/ecstore/src/cache_value/metacache_set.rs index f174f502f..c19bd314e 100644 --- a/crates/ecstore/src/cache_value/metacache_set.rs +++ b/crates/ecstore/src/cache_value/metacache_set.rs @@ -15,7 +15,6 @@ use crate::disk::disk_store::get_drive_walkdir_stall_timeout; use crate::disk::error::DiskError; 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::{ @@ -28,8 +27,8 @@ use std::{ time::Duration, }; use tokio::io::{AsyncRead, AsyncWrite}; -use tokio::spawn; use tokio::sync::Mutex as TokioMutex; +use tokio::task::JoinSet; use tokio::time::timeout; use tokio_util::sync::CancellationToken; use tracing::{debug, error, warn}; @@ -75,13 +74,22 @@ enum PeekOutcome { Ready(Option), Error(rustfs_filemeta::Error), TimedOut, + Cancelled, } -async fn peek_with_timeout(reader: &mut MetacacheReader, timeout_duration: Duration) -> PeekOutcome { - match timeout(timeout_duration, reader.peek()).await { - Ok(Ok(entry)) => PeekOutcome::Ready(entry), - Ok(Err(err)) => PeekOutcome::Error(err), - Err(_) => PeekOutcome::TimedOut, +async fn peek_with_timeout( + cancel: &CancellationToken, + reader: &mut MetacacheReader, + timeout_duration: Duration, +) -> PeekOutcome { + tokio::select! { + biased; + _ = cancel.cancelled() => PeekOutcome::Cancelled, + result = timeout(timeout_duration, reader.peek()) => match result { + Ok(Ok(entry)) => PeekOutcome::Ready(entry), + Ok(Err(err)) => PeekOutcome::Error(err), + Err(_) => PeekOutcome::TimedOut, + }, } } @@ -159,12 +167,26 @@ pub(crate) enum TestReaderBehavior { Eof, Entries(Vec), Stall, + StallWithSignals { + entered: Arc, + stopped: Arc, + }, IgnoreCancel, ProducerError(DiskError), PrimaryErrorThenFallback(DiskError), PartialThenTimeout(Vec), } +#[cfg(test)] +struct NotifyOnDrop(Arc); + +#[cfg(test)] +impl Drop for NotifyOnDrop { + fn drop(&mut self) { + self.0.notify_one(); + } +} + #[derive(Default)] pub struct ListPathRawOptions { pub disks: Vec>, @@ -224,6 +246,8 @@ impl Clone for ListPathRawOptions { } pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> disk::error::Result<()> { + let rx = rx.child_token(); + let _cancel_guard = rx.clone().drop_guard(); list_path_raw_inner(rx, opts, None).await } @@ -232,6 +256,8 @@ pub(crate) async fn list_path_raw_with_claim_tracker( opts: ListPathRawOptions, claim_tracker: FallbackClaimTracker, ) -> disk::error::Result<()> { + let rx = rx.child_token(); + let _cancel_guard = rx.clone().drop_guard(); list_path_raw_inner(rx, opts, Some(claim_tracker)).await } @@ -250,7 +276,7 @@ async fn list_path_raw_inner( let log_bucket = opts.bucket.clone(); let log_path = opts.path.clone(); - let mut jobs: Vec>> = Vec::new(); + let mut jobs = JoinSet::new(); let mut readers = Vec::with_capacity(opts.disks.len()); let fds = Arc::new(TokioMutex::new(opts.fallback_disks.iter().flatten().cloned().collect::>())); #[cfg(test)] @@ -273,7 +299,7 @@ async fn list_path_raw_inner( let producer_errs_clone = producer_errs.clone(); let (rd, wr) = tokio::io::duplex(64); readers.push(MetacacheReader::new(rd)); - jobs.push(spawn(async move { + jobs.spawn(async move { let mut wr = PublishedBytesWriter::new(wr); #[cfg(test)] let test_primary_error = if let Some(behavior) = opts_clone.test_reader_behaviors.get(disk_idx).cloned() { @@ -290,6 +316,12 @@ async fn list_path_raw_inner( cancel_rx_clone.cancelled().await; return Ok(()); } + TestReaderBehavior::StallWithSignals { entered, stopped } => { + let _drop_signal = NotifyOnDrop(stopped); + entered.notify_one(); + let _held_writer = &mut wr; + std::future::pending::>().await + } TestReaderBehavior::IgnoreCancel => { let _held_writer = wr; std::future::pending::<()>().await; @@ -427,7 +459,9 @@ async fn list_path_raw_inner( last_err = Some(err); continue; } - TestReaderBehavior::Stall | TestReaderBehavior::IgnoreCancel => { + TestReaderBehavior::Stall + | TestReaderBehavior::StallWithSignals { .. } + | TestReaderBehavior::IgnoreCancel => { last_err = Some(DiskError::Timeout); continue; } @@ -566,10 +600,12 @@ async fn list_path_raw_inner( // warn!("list_path_raw: while need_fallback done"); Ok(()) - })); + }); } - let revjob = spawn(async move { + let revjob_rx = rx.clone(); + let mut revjobs = JoinSet::new(); + revjobs.spawn(async move { let peek_timeout = opts .walkdir_stall_timeout .or({ @@ -597,7 +633,7 @@ async fn list_path_raw_inner( // opts.bucket, opts.path, ¤t.name // ); - if rx.is_cancelled() { + if revjob_rx.is_cancelled() { return Ok(()); } @@ -618,7 +654,7 @@ async fn list_path_raw_inner( let entry = if let Some(entry) = pending_entries[i].take() { entry } else { - match peek_with_timeout(r, peek_timeout).await { + match peek_with_timeout(&revjob_rx, r, peek_timeout).await { PeekOutcome::Ready(res) => { if let Some(entry) = res { // info!("read entry disk: {}, name: {}", i, entry.name); @@ -703,6 +739,7 @@ async fn list_path_raw_inner( *r = MetacacheReader::new(detached_rd); continue; } + PeekOutcome::Cancelled => return Ok(()), } }; @@ -757,7 +794,11 @@ async fn list_path_raw_inner( if has_err > 0 && has_err > opts.disks.len() - opts.min_disks { if let Some(finished_fn) = opts.finished.as_ref() { - finished_fn(&errs).await; + tokio::select! { + biased; + _ = revjob_rx.cancelled() => return Ok(()), + _ = finished_fn(&errs) => {} + } } let mut combined_err = Vec::new(); errs.iter().zip(opts.disks.iter()).for_each(|(err, disk)| match (err, disk) { @@ -789,7 +830,11 @@ async fn list_path_raw_inner( if has_err > 0 && let Some(finished_fn) = opts.finished.as_ref() { - finished_fn(&errs).await; + tokio::select! { + biased; + _ = revjob_rx.cancelled() => return Ok(()), + _ = finished_fn(&errs) => {} + } } // Tolerated reader failures, including timeouts, must not turn a // quorum EOF into a failed listing. The quorum failure branch above @@ -805,7 +850,11 @@ async fn list_path_raw_inner( if let Some(agreed_fn) = opts.agreed.as_ref() { // warn!("list_path_raw: agreed_fn start, current: {:?}", ¤t.name); - agreed_fn(current).await; + tokio::select! { + biased; + _ = revjob_rx.cancelled() => return Ok(()), + _ = agreed_fn(current) => {} + } // warn!("list_path_raw: agreed_fn done"); } @@ -821,14 +870,21 @@ async fn list_path_raw_inner( } if let Some(partial_fn) = opts.partial.as_ref() { - partial_fn(MetaCacheEntries(top_entries), &errs).await; + tokio::select! { + biased; + _ = revjob_rx.cancelled() => return Ok(()), + _ = partial_fn(MetaCacheEntries(top_entries), &errs) => {} + } } } Ok(()) }); let merge_started = std::time::Instant::now(); - if let Err(err) = revjob.await.map_err(std::io::Error::other)? { + let Some(revjob) = revjobs.join_next().await else { + return Err(DiskError::Unexpected); + }; + if let Err(err) = revjob.map_err(std::io::Error::other)? { rustfs_io_metrics::record_stage_duration("metacache_merge_failed", merge_started.elapsed().as_secs_f64() * 1000.0); if is_missing_path_error(&err) { debug!( @@ -854,12 +910,15 @@ async fn list_path_raw_inner( ); } cancel_rx.cancel(); - for job in jobs { - job.abort(); - } + jobs.abort_all(); return Err(err); } + if rx.is_cancelled() { + cancel_rx.cancel(); + jobs.abort_all(); + return Ok(()); + } rustfs_io_metrics::record_stage_duration("metacache_merge", merge_started.elapsed().as_secs_f64() * 1000.0); // The merge consumer can finish successfully before every producer finishes @@ -867,16 +926,9 @@ async fn list_path_raw_inner( // or after the requested listing limit is satisfied). Cancel remaining walk // jobs before aborting them so list calls do not wait for slow remote streams. cancel_rx.cancel(); - - for job in jobs.iter() { - if !job.is_finished() { - job.abort(); - } - } - - let results = join_all(jobs).await; + jobs.abort_all(); let mut job_errs = Vec::new(); - for result in results { + while let Some(result) = jobs.join_next().await { match result { Ok(Ok(())) => {} Ok(Err(err)) => { @@ -1285,6 +1337,226 @@ mod tests { .expect("external cancellation should stop listing without a synthetic canceled error"); } + #[tokio::test] + async fn list_path_raw_cancels_a_blocked_peek_without_waiting_for_the_stall_timeout() { + let cancel = CancellationToken::new(); + let entered = Arc::new(tokio::sync::Notify::new()); + let stopped = Arc::new(tokio::sync::Notify::new()); + let list = tokio::spawn(list_path_raw( + cancel.clone(), + ListPathRawOptions { + disks: vec![None], + min_disks: 1, + test_reader_behaviors: vec![TestReaderBehavior::StallWithSignals { + entered: entered.clone(), + stopped: stopped.clone(), + }], + peek_timeout: Some(Duration::from_secs(60)), + ..Default::default() + }, + )); + + timeout(Duration::from_secs(1), entered.notified()) + .await + .expect("producer should block before cancellation"); + cancel.cancel(); + + let result = timeout(Duration::from_secs(1), list) + .await + .expect("cancellation should interrupt the blocked metacache peek") + .expect("listing task should not panic"); + assert!(result.is_ok(), "cancellation should stop the listing cleanly, got {result:?}"); + timeout(Duration::from_secs(1), stopped.notified()) + .await + .expect("cancellation should drop the blocked producer"); + } + + #[tokio::test] + async fn aborting_list_path_raw_drops_blocked_producers() { + let entered = Arc::new(tokio::sync::Notify::new()); + let stopped = Arc::new(tokio::sync::Notify::new()); + let list = tokio::spawn(list_path_raw( + CancellationToken::new(), + ListPathRawOptions { + disks: vec![None], + min_disks: 1, + test_reader_behaviors: vec![TestReaderBehavior::StallWithSignals { + entered: entered.clone(), + stopped: stopped.clone(), + }], + peek_timeout: Some(Duration::from_secs(60)), + ..Default::default() + }, + )); + + timeout(Duration::from_secs(1), entered.notified()) + .await + .expect("producer should block before aborting the listing"); + list.abort(); + assert!(list.await.expect_err("listing task should be cancelled").is_cancelled()); + timeout(Duration::from_secs(1), stopped.notified()) + .await + .expect("aborting the listing should drop the blocked producer"); + } + + #[tokio::test] + async fn list_path_raw_cancels_a_blocked_partial_callback() { + let cancel = CancellationToken::new(); + let callback_started = Arc::new(tokio::sync::Notify::new()); + let callback_started_clone = callback_started.clone(); + let entry = MetaCacheEntry { + name: "bucket/object-a".to_string(), + metadata: vec![1], + ..Default::default() + }; + let later_entry = MetaCacheEntry { + name: "bucket/object-b".to_string(), + metadata: vec![1], + ..Default::default() + }; + let list = tokio::spawn(list_path_raw( + cancel.clone(), + ListPathRawOptions { + disks: vec![None, None], + min_disks: 2, + test_reader_behaviors: vec![ + TestReaderBehavior::Entries(vec![entry]), + TestReaderBehavior::Entries(vec![later_entry]), + ], + partial: Some(Box::new(move |_, _| { + let callback_started = callback_started_clone.clone(); + Box::pin(async move { + callback_started.notify_one(); + std::future::pending::<()>().await; + }) + })), + ..Default::default() + }, + )); + + timeout(Duration::from_secs(1), callback_started.notified()) + .await + .expect("partial callback should block before cancellation"); + cancel.cancel(); + + let result = timeout(Duration::from_secs(1), list) + .await + .expect("cancellation should interrupt the blocked partial callback") + .expect("listing task should not panic"); + assert!(result.is_ok(), "cancellation should stop the listing cleanly, got {result:?}"); + } + + #[tokio::test] + async fn list_path_raw_cancels_a_blocked_agreed_callback() { + let cancel = CancellationToken::new(); + let callback_started = Arc::new(tokio::sync::Notify::new()); + let callback_started_clone = callback_started.clone(); + let entry = MetaCacheEntry { + name: "bucket/object".to_string(), + metadata: vec![1], + ..Default::default() + }; + let list = tokio::spawn(list_path_raw( + cancel.clone(), + ListPathRawOptions { + disks: vec![None], + min_disks: 1, + test_reader_behaviors: vec![TestReaderBehavior::Entries(vec![entry])], + agreed: Some(Box::new(move |_| { + let callback_started = callback_started_clone.clone(); + Box::pin(async move { + callback_started.notify_one(); + std::future::pending::<()>().await; + }) + })), + ..Default::default() + }, + )); + + timeout(Duration::from_secs(1), callback_started.notified()) + .await + .expect("agreed callback should block before cancellation"); + cancel.cancel(); + + let result = timeout(Duration::from_secs(1), list) + .await + .expect("cancellation should interrupt the blocked agreed callback") + .expect("listing task should not panic"); + assert!(result.is_ok(), "cancellation should stop the listing cleanly, got {result:?}"); + } + + #[tokio::test] + async fn list_path_raw_cancels_a_blocked_finished_callback_after_quorum_failure() { + let cancel = CancellationToken::new(); + let callback_started = Arc::new(tokio::sync::Notify::new()); + let callback_started_clone = callback_started.clone(); + let list = tokio::spawn(list_path_raw( + cancel.clone(), + ListPathRawOptions { + disks: vec![None], + min_disks: 1, + test_reader_behaviors: vec![TestReaderBehavior::ProducerError(DiskError::FileAccessDenied)], + finished: Some(Box::new(move |_| { + let callback_started = callback_started_clone.clone(); + Box::pin(async move { + callback_started.notify_one(); + std::future::pending::<()>().await; + }) + })), + ..Default::default() + }, + )); + + timeout(Duration::from_secs(1), callback_started.notified()) + .await + .expect("finished callback should block before cancellation"); + cancel.cancel(); + + let result = timeout(Duration::from_secs(1), list) + .await + .expect("cancellation should interrupt the blocked finished callback") + .expect("listing task should not panic"); + assert!(result.is_ok(), "cancellation should stop the listing cleanly, got {result:?}"); + } + + #[tokio::test] + async fn list_path_raw_cancels_a_blocked_finished_callback_after_tolerated_error() { + let cancel = CancellationToken::new(); + let callback_started = Arc::new(tokio::sync::Notify::new()); + let callback_started_clone = callback_started.clone(); + let list = tokio::spawn(list_path_raw( + cancel.clone(), + ListPathRawOptions { + disks: vec![None, None, None], + min_disks: 2, + test_reader_behaviors: vec![ + TestReaderBehavior::Eof, + TestReaderBehavior::Eof, + TestReaderBehavior::ProducerError(DiskError::FileAccessDenied), + ], + finished: Some(Box::new(move |_| { + let callback_started = callback_started_clone.clone(); + Box::pin(async move { + callback_started.notify_one(); + std::future::pending::<()>().await; + }) + })), + ..Default::default() + }, + )); + + timeout(Duration::from_secs(1), callback_started.notified()) + .await + .expect("finished callback should block before cancellation"); + cancel.cancel(); + + let result = timeout(Duration::from_secs(1), list) + .await + .expect("cancellation should interrupt the blocked finished callback") + .expect("listing task should not panic"); + assert!(result.is_ok(), "cancellation should stop the listing cleanly, got {result:?}"); + } + #[tokio::test] async fn list_path_raw_ignores_closed_output_stream_after_quorum_eof() { list_path_raw( @@ -1491,15 +1763,50 @@ mod tests { async fn peek_with_timeout_times_out_on_silent_reader() { let (_writer, reader) = tokio::io::duplex(64); let mut reader = MetacacheReader::new(reader); + let cancel = CancellationToken::new(); - let outcome = peek_with_timeout(&mut reader, Duration::from_millis(20)).await; + let outcome = peek_with_timeout(&cancel, &mut reader, Duration::from_millis(20)).await; assert!(matches!(outcome, PeekOutcome::TimedOut)); } + #[tokio::test] + async fn peek_with_timeout_stops_immediately_when_cancelled() { + struct PendingReader(Arc); + + impl AsyncRead for PendingReader { + fn poll_read( + self: std::pin::Pin<&mut Self>, + _cx: &mut std::task::Context<'_>, + _buf: &mut tokio::io::ReadBuf<'_>, + ) -> std::task::Poll> { + self.get_mut().0.notify_one(); + std::task::Poll::Pending + } + } + + let entered = Arc::new(tokio::sync::Notify::new()); + let mut reader = MetacacheReader::new(PendingReader(entered.clone())); + let cancel = CancellationToken::new(); + let wait_cancel = cancel.clone(); + + let wait = tokio::spawn(async move { peek_with_timeout(&wait_cancel, &mut reader, Duration::from_secs(60)).await }); + timeout(Duration::from_secs(1), entered.notified()) + .await + .expect("peek should block before cancellation"); + cancel.cancel(); + + let outcome = timeout(Duration::from_secs(1), wait) + .await + .expect("cancellation should interrupt the pending peek") + .expect("peek task should not panic"); + assert!(matches!(outcome, PeekOutcome::Cancelled)); + } + #[tokio::test] async fn peek_with_timeout_reads_entry_before_deadline() { let (reader, writer) = tokio::io::duplex(256); let mut metacache_reader = MetacacheReader::new(reader); + let cancel = CancellationToken::new(); tokio::spawn(async move { let mut writer = MetacacheWriter::new(writer); @@ -1513,7 +1820,7 @@ mod tests { writer.close().await.expect("writer should close"); }); - let outcome = peek_with_timeout(&mut metacache_reader, Duration::from_secs(1)).await; + let outcome = peek_with_timeout(&cancel, &mut metacache_reader, Duration::from_secs(1)).await; match outcome { PeekOutcome::Ready(Some(entry)) => assert_eq!(entry.name, "bucket/object"), other => panic!("expected ready entry, got {other:?}"), diff --git a/crates/ecstore/src/store/list_objects.rs b/crates/ecstore/src/store/list_objects.rs index 2d2780f39..93e986244 100644 --- a/crates/ecstore/src/store/list_objects.rs +++ b/crates/ecstore/src/store/list_objects.rs @@ -60,6 +60,7 @@ use tokio::io::duplex; use tokio::sync::broadcast::{self}; use tokio::sync::mpsc::{self, Receiver, Sender}; use tokio::sync::{OnceCell, RwLock}; +use tokio::task::JoinSet; use tokio_util::sync::CancellationToken; use tracing::{Instrument, debug, error, info, warn}; use uuid::Uuid; @@ -2660,7 +2661,8 @@ async fn read_fallback_listing_disk( ) -> Vec { let (rd, mut wr) = duplex(64); let walk_endpoint = endpoint.clone(); - let walk_job = tokio::spawn(async move { + let mut walk_jobs = JoinSet::new(); + walk_jobs.spawn(async move { disk.walk_dir( WalkDirOptions { bucket: options.bucket, @@ -2716,10 +2718,10 @@ async fn read_fallback_listing_disk( } drop(reader); - match walk_job.await { - Ok(Ok(())) if stream_ok => entries, - Ok(Ok(())) => Vec::new(), - Ok(Err(err)) => { + match walk_jobs.join_next().await { + Some(Ok(Ok(()))) if stream_ok => entries, + Some(Ok(Ok(()))) => Vec::new(), + Some(Ok(Err(err))) => { debug!( endpoint = %walk_endpoint, error = ?err, @@ -2727,7 +2729,7 @@ async fn read_fallback_listing_disk( ); Vec::new() } - Err(err) => { + Some(Err(err)) => { debug!( endpoint = %walk_endpoint, error = ?err, @@ -2735,6 +2737,7 @@ async fn read_fallback_listing_disk( ); Vec::new() } + None => Vec::new(), } } @@ -3986,6 +3989,7 @@ impl ECStore { // cancel channel let cancel = CancellationToken::new(); + let _cancel_guard = cancel.clone().drop_guard(); let (err_tx, mut err_rx) = broadcast::channel::>(1); @@ -5226,6 +5230,7 @@ impl Sets { let log_context = ListPathLogContext::from_options(&o); let cancel = CancellationToken::new(); + let _cancel_guard = cancel.clone().drop_guard(); let (err_tx, mut err_rx) = broadcast::channel::>(1); let (sender, recv) = mpsc::channel(o.limit as usize); @@ -6180,6 +6185,7 @@ impl SetDisks { let log_context = ListPathLogContext::from_options(&o); let cancel = CancellationToken::new(); + let _cancel_guard = cancel.clone().drop_guard(); let (err_tx, mut err_rx) = broadcast::channel::>(1); let (sender, recv) = mpsc::channel(o.limit as usize);