From 513bc8249950fae8535691e2bd9bb9290301e187 Mon Sep 17 00:00:00 2001 From: overtrue Date: Fri, 14 Aug 2026 01:57:07 +0800 Subject: [PATCH] chore(ecstore): remove the pool-level ListObjects pagination copy MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The ListObjects pagination pipeline existed in three near-copies in one file; production listing never reaches the Sets copy, which ECStore bypasses by expanding straight to per-set disks. This removes it: impl ListOperations for Sets (61 lines of pure forwarding in core/sets.rs) and the impl Sets pagination block (826 lines of inner_list_objects_v2 / list_objects_generic / inner_list_object_versions / list_path / list_merged / walk_internal in store/list_objects.rs). Two preconditions verified before deleting rather than taken on faith: the architecture guard pins only set_disks_implements_storage_list_operations_contract, so nothing requires the Sets trait impl; and the four Sets pagination methods had no cross-file caller besides that trait impl. The single test consumer moves to the surviving pipeline instead of being deleted: writes still go through the pool, and the listing assertion now targets the set-level implementation. It is renamed accordingly so the name still describes what it covers. The logging guardrail's TRACE-only requirement for Sets::list_objects_v2 retires in the same diff — the wrapper it pinned no longer exists. The ECStore and SetDisks entries are untouched. The SetDisks copy stays for now: its trait impl is guard-pinned, so replacing the duplicate pipeline behind it needs the generic helper the issue schedules for post-1.0. Verification: cargo nextest run -p rustfs-ecstore 4020 passed; check_architecture_migration_rules.sh and check_logging_guardrails.sh pass; clippy --lib --tests -D warnings clean; make pre-commit green. Ref rustfs/backlog#1821 (PR1). --- crates/ecstore/src/core/sets.rs | 73 +- crates/ecstore/src/store/list_objects.rs | 827 ----------------------- scripts/check_logging_guardrails.sh | 4 +- 3 files changed, 11 insertions(+), 893 deletions(-) diff --git a/crates/ecstore/src/core/sets.rs b/crates/ecstore/src/core/sets.rs index acb8b53b9..bb1026d28 100644 --- a/crates/ecstore/src/core/sets.rs +++ b/crates/ecstore/src/core/sets.rs @@ -59,7 +59,6 @@ use std::{ use tokio::sync::RwLock; use tokio::sync::broadcast::{Receiver, Sender}; use tokio::time::Duration; -use tokio_util::sync::CancellationToken; use tracing::warn; use tracing::{error, info}; use uuid::Uuid; @@ -756,67 +755,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for Sets { } } -#[async_trait::async_trait] -impl crate::storage_api_contracts::list::ListOperations for Sets { - type Error = Error; - type ListObjectsV2Info = ListObjectsV2Info; - type ListObjectVersionsInfo = ListObjectVersionsInfo; - type ObjectInfoOrErr = ObjectInfoOrErr; - type WalkOptions = WalkOptions; - type WalkCancellation = CancellationToken; - type WalkResultSender = tokio::sync::mpsc::Sender; - - #[tracing::instrument(level = "trace", skip(self))] - async fn list_objects_v2( - self: Arc, - bucket: &str, - prefix: &str, - continuation_token: Option, - delimiter: Option, - max_keys: i32, - fetch_owner: bool, - start_after: Option, - incl_deleted: bool, - ) -> Result { - self.inner_list_objects_v2( - bucket, - prefix, - continuation_token, - delimiter, - max_keys, - fetch_owner, - start_after, - incl_deleted, - ) - .await - } - - #[tracing::instrument(skip(self))] - async fn list_object_versions( - self: Arc, - bucket: &str, - prefix: &str, - marker: Option, - version_marker: Option, - delimiter: Option, - max_keys: i32, - ) -> Result { - self.inner_list_object_versions(bucket, prefix, marker, version_marker, delimiter, max_keys) - .await - } - - async fn walk( - self: Arc, - rx: CancellationToken, - bucket: &str, - prefix: &str, - result: tokio::sync::mpsc::Sender, - opts: WalkOptions, - ) -> Result<()> { - self.walk_internal(rx, bucket, prefix, result, opts).await - } -} - #[async_trait::async_trait] impl crate::storage_api_contracts::multipart::MultipartOperations for Sets { type Error = Error; @@ -1903,7 +1841,7 @@ mod tests { #[tokio::test(flavor = "multi_thread")] #[serial] - async fn sets_list_objects_v2_lists_objects_within_the_pool() { + async fn set_level_list_objects_v2_lists_objects_written_through_the_pool() { let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::Erasure).await; let format = FormatV3::new(1, 2); let mut endpoints = Vec::new(); @@ -1985,11 +1923,16 @@ mod tests { .await .expect("object should be written"); - let result = sets + // Listing asserts against the set-level pipeline: the pool-level + // ListOperations impl for `Sets` was a pure forwarder over a duplicate + // pagination block that production never reached, and both were removed + // (backlog#1821). Writes still go through the pool, so this keeps + // covering the write-then-list round trip end to end. + let result = sets.disk_set[0] .clone() .list_objects_v2(&bucket, "", None, None, 1000, false, None, false) .await - .expect("pool-level listing should succeed"); + .expect("set-level listing should succeed"); assert_eq!(result.objects.len(), 1); assert_eq!(result.objects[0].name, object); diff --git a/crates/ecstore/src/store/list_objects.rs b/crates/ecstore/src/store/list_objects.rs index 2e8ee9968..df4cb3980 100644 --- a/crates/ecstore/src/store/list_objects.rs +++ b/crates/ecstore/src/store/list_objects.rs @@ -16,7 +16,6 @@ use crate::bucket::metadata_sys::get_versioning_config; use crate::bucket::utils::check_list_objs_args; use crate::bucket::versioning::VersioningApi; use crate::cache_value::metacache_set::{FallbackClaimTracker, ListPathRawOptions, list_path_raw_with_claim_tracker}; -use crate::core::sets::Sets; use crate::disk::error::DiskError; use crate::disk::{DiskAPI, DiskInfo, DiskStore, RUSTFS_META_BUCKET, WalkDirOptions}; use crate::error::{ @@ -5030,832 +5029,6 @@ async fn merge_entry_channels( Ok(()) } -impl Sets { - #[allow(clippy::too_many_arguments)] - pub async fn inner_list_objects_v2( - self: Arc, - bucket: &str, - prefix: &str, - continuation_token: Option, - delimiter: Option, - max_keys: i32, - _fetch_owner: bool, - start_after: Option, - incl_deleted: bool, - ) -> Result { - let marker = if continuation_token.is_none() { - start_after - } else { - continuation_token.clone() - }; - - let loi = self - .list_objects_generic(bucket, prefix, marker, delimiter, max_keys, incl_deleted) - .await?; - Ok(ListObjectsV2Info { - is_truncated: loi.is_truncated, - continuation_token, - next_continuation_token: loi.next_marker, - objects: loi.objects, - prefixes: loi.prefixes, - }) - } - - pub async fn list_objects_generic( - self: Arc, - bucket: &str, - prefix: &str, - marker: Option, - delimiter: Option, - max_keys: i32, - incl_deleted: bool, - ) -> Result { - let max_keys = normalize_max_keys(max_keys); - let effective_max_keys = if max_keys <= 0 { 0 } else { max_keys_plus_one(max_keys, true) }; - let mut opts = ListPathOptions { - bucket: bucket.to_owned(), - prefix: prefix.to_owned(), - separator: delimiter.clone(), - limit: effective_max_keys, - marker, - incl_deleted, - ask_disks: list_objects_quorum_from_env(), - ..Default::default() - }; - // Strip the `[rustfs_cache:...]` cursor tag before any name comparison - // (notably `forward_past`) — see backlog#1047. - opts.parse_marker(); - - if !opts.prefix.is_empty() && max_keys == 1 && opts.marker.is_none() && !incl_deleted { - match self - .get_object_info( - &opts.bucket, - &opts.prefix, - &ObjectOptions { - no_lock: true, - ..Default::default() - }, - ) - .await - { - Ok(res) if !res.delete_marker => { - return Ok(ListObjectsInfo { - objects: vec![res], - ..Default::default() - }); - } - Err(err) if is_err_bucket_not_found(&err) => { - return Err(err); - } - _ => {} - }; - } - - let mut list_result = self - .list_path(&opts) - .await - .unwrap_or_else(|err| MetaCacheEntriesSortedResult { - err: Some(err.into()), - ..Default::default() - }); - let next_cache_id = list_result.entries.as_ref().and_then(|entries| entries.list_id.clone()); - - let disk_has_more = list_result.err.is_none(); - - if let Some(err) = list_result.err.take() - && err != rustfs_filemeta::Error::Unexpected - { - return Err(to_object_err(err.into(), vec![bucket, prefix])); - } - - if let Some(result) = list_result.entries.as_mut() { - result.forward_past(opts.marker); - } - - // Last RAW scanned key, captured before folding (ECA-03 / #944). - let last_scanned_key = last_scanned_entry_name(list_result.entries.as_ref()); - - let get_objects = ObjectInfo::from_meta_cache_entries_sorted_infos( - &list_result.entries.unwrap_or_default(), - bucket, - prefix, - delimiter.clone(), - ) - .await; - - let (objects, prefixes, is_truncated, next_marker, next_version_idmarker) = list_objects_paginate( - get_objects, - &delimiter, - max_keys, - disk_has_more, - next_cache_id.as_deref(), - false, - last_scanned_key.as_deref(), - ); - let _ = next_version_idmarker; - - Ok(ListObjectsInfo { - is_truncated, - next_marker, - objects, - prefixes, - }) - } - - pub async fn inner_list_object_versions( - self: Arc, - bucket: &str, - prefix: &str, - marker: Option, - version_marker: Option, - delimiter: Option, - max_keys: i32, - ) -> Result { - let max_keys = normalize_max_keys(max_keys); - if marker.is_none() && version_marker.is_some() { - return Err(StorageError::NotImplemented); - } - - let has_version_marker = version_marker.is_some(); - let version_marker = if let Some(marker) = version_marker { - Some(parse_version_marker(marker)?) - } else { - None - }; - - let effective_max_keys = if max_keys <= 0 { 0 } else { max_keys_plus_one(max_keys, true) }; - let mut opts = ListPathOptions { - bucket: bucket.to_owned(), - prefix: prefix.to_owned(), - separator: delimiter.clone(), - limit: effective_max_keys, - marker, - incl_deleted: true, - ask_disks: list_objects_quorum_from_env(), - versioned: true, - include_marker: has_version_marker, - ..Default::default() - }; - // Strip the `[rustfs_cache:...]` cursor tag before any name comparison - // (notably `forward_past`) — see backlog#1047. - opts.parse_marker(); - - let mut list_result = self - .list_path(&opts) - .await - .unwrap_or_else(|err| MetaCacheEntriesSortedResult { - err: Some(err.into()), - ..Default::default() - }); - let next_cache_id = list_result.entries.as_ref().and_then(|entries| entries.list_id.clone()); - - let disk_has_more = list_result.err.is_none(); - - if let Some(err) = list_result.err.take() - && err != rustfs_filemeta::Error::Unexpected - { - return Err(to_object_err(err.into(), vec![bucket, prefix])); - } - - if let Some(result) = list_result.entries.as_mut() - && !has_version_marker - { - result.forward_past(opts.marker.clone()); - } - - let version_marker = version_marker_for_entries(list_result.entries.as_ref(), opts.marker.as_deref(), version_marker); - - // Last RAW scanned key, captured before folding (ECA-03 / #944). - let last_scanned_key = last_scanned_entry_name(list_result.entries.as_ref()); - - let get_objects = ObjectInfo::from_meta_cache_entries_sorted_versions( - &list_result.entries.unwrap_or_default(), - bucket, - prefix, - delimiter.clone(), - version_marker, - ) - .await; - - let (objects, prefixes, is_truncated, next_marker, next_version_idmarker) = list_objects_paginate( - get_objects, - &delimiter, - max_keys, - disk_has_more, - next_cache_id.as_deref(), - true, - last_scanned_key.as_deref(), - ); - - Ok(ListObjectVersionsInfo { - is_truncated, - next_marker, - next_version_idmarker, - objects, - prefixes, - }) - } - - pub async fn list_path(self: Arc, o: &ListPathOptions) -> Result { - check_list_objs_args(&o.bucket, &o.prefix, &o.marker)?; - let list_path_started = std::time::Instant::now(); - - let mut o = o.clone(); - o.marker = o.marker.filter(|v| v >= &o.prefix); - - if let Some(marker) = &o.marker - && !o.prefix.is_empty() - && !marker.starts_with(&o.prefix) - { - return Err(Error::Unexpected); - } - - if o.limit == 0 { - return Err(Error::Unexpected); - } - - if o.prefix.starts_with(SLASH_SEPARATOR) { - return Err(Error::Unexpected); - } - - let slash_separator = Some(SLASH_SEPARATOR.to_owned()); - o.include_directories = o.separator == slash_separator; - - if (o.separator == slash_separator || o.separator.is_none()) && !o.recursive { - o.recursive = o.separator != slash_separator; - o.separator = slash_separator; - } else { - o.recursive = true; - } - - o.parse_marker(); - - if o.base_dir.is_empty() { - o.base_dir = base_dir_from_prefix(&o.prefix); - } - - o.transient = o.transient || is_reserved_or_invalid_bucket(&o.bucket, false); - o.set_filter(); - if o.transient { - o.create = false; - } - 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); - - let sets = self.clone(); - let opts = o.clone(); - let cancel_rx1 = cancel.clone(); - let cancel_rx1_for_err = cancel_rx1.clone(); - let err_tx1 = err_tx.clone(); - let job1_context = log_context.clone(); - let job1 = tokio::spawn( - async move { - let mut opts = opts; - opts.stop_disk_at_limit = true; - if let Err(err) = sets.list_merged(cancel_rx1, opts, sender).await - && !cancel_rx1_for_err.is_cancelled() - { - log_list_path_worker_error("sets", "list_merged", &job1_context, &err); - let _ = err_tx1.send(Arc::new(err)); - } - } - .instrument(tracing::Span::current()), - ); - - let cancel_rx2 = cancel.clone(); - let (result_tx, mut result_rx) = mpsc::channel(1); - let err_tx2 = err_tx.clone(); - let opts = o.clone(); - let job2_context = log_context.clone(); - let job2 = tokio::spawn( - async move { - match gather_results(cancel_rx2, opts, recv, result_tx).await { - Ok(GatherResultsState::LimitReached) => cancel.cancel(), - Ok(GatherResultsState::InputClosed) => {} - // Consumer disconnect (e.g. client cancelled the request) - // is a benign completion: no error log, no err_tx send. - // The explicit cancel is idempotent and avoids relying on - // the wrapper's drop-guard ordering to stop the producer. - Ok(GatherResultsState::ConsumerGone) => cancel.cancel(), - // Invariant (rustfs/backlog#1306): gather_results maps a - // consumer disconnect to Ok(ConsumerGone) and has no - // fallible pre-send path, so this arm is currently - // unreachable. It is kept as a guard: it stays correct only - // while a *real* gather_results error would still surface as - // an error here. Real producer/listing errors take the - // separate job1 -> err_tx -> err_rx path below, unaffected. - Err(err) => { - log_list_path_worker_error("sets", "gather_results", &job2_context, &err); - let _ = err_tx2.send(Arc::new(err)); - cancel.cancel(); - } - } - } - .instrument(tracing::Span::current()), - ); - - let mut result = tokio::select! { - res = err_rx.recv() => { - match res { - Ok(err) => { - log_list_path_worker_error("sets", "worker_error", &log_context, err.as_ref()); - MetaCacheEntriesSortedResult { entries: None, err: Some(err.as_ref().clone().into()) } - }, - Err(err) => { - log_list_path_worker_error("sets", "error_channel_closed", &log_context, &err); - MetaCacheEntriesSortedResult { entries: None, err: Some(rustfs_filemeta::Error::other(err)) } - }, - } - } - Some(result) = result_rx.recv() => result, - }; - - join_all(vec![job1, job2]).await; - - if let Ok(err) = err_rx.try_recv() { - log_list_path_worker_error("sets", "trailing_worker_error", &log_context, err.as_ref()); - result.err = Some(err.as_ref().clone().into()); - } - - if result.err.is_some() { - log_list_path_finished("sets", &log_context, list_path_started.elapsed().as_secs_f64() * 1000.0, 0, true); - return Ok(result); - } - - if let Some(entries) = result.entries.as_mut() { - entries.reuse = true; - let truncated = !entries.entries().is_empty() || result.err.is_none(); - entries.o.0.truncate(o.limit as usize); - if !o.transient && truncated { - entries.list_id = if let Some(id) = o.id { - Some(id) - } else { - Some(Uuid::new_v4().to_string()) - }; - } - - if !truncated { - result.err = Some(Error::Unexpected.into()); - } - } - - log_list_path_finished( - "sets", - &log_context, - list_path_started.elapsed().as_secs_f64() * 1000.0, - result - .entries - .as_ref() - .map(|entries| entries.entries().len()) - .unwrap_or_default(), - result.err.is_some(), - ); - - Ok(result) - } - - async fn list_merged( - &self, - rx: CancellationToken, - opts: ListPathOptions, - sender: Sender, - ) -> Result> { - let merge_started = std::time::Instant::now(); - - debug!( - list_path_limit = opts.limit, - recursive = opts.recursive, - set_count = self.disk_set.len(), - requested_marker = %opts.marker.as_deref().unwrap_or(""), - "sets list_merged started" - ); - - let mut futures = Vec::new(); - let mut inputs = Vec::new(); - - for set in &self.disk_set { - let (send, recv) = list_merged_entry_channel(); - inputs.push(recv); - let opts = opts.clone(); - let rx_clone = rx.clone(); - let set = set.clone(); - futures.push(async move { set.list_path(rx_clone, opts, send).await }); - } - - tokio::spawn( - async move { - if let Err(err) = merge_entry_channels(rx, inputs, sender.clone(), 1).await { - error!("merge_entry_channels err {:?}", err); - } - } - .instrument(tracing::Span::current()), - ); - - let results = join_all(futures).await; - let mut all_at_eof = true; - let mut errs = Vec::new(); - for result in results { - if let Err(err) = result { - all_at_eof = false; - errs.push(Some(err)); - } else { - errs.push(None); - } - } - - if is_all_not_found(&errs) { - return Ok(Vec::new()); - } - - for err in &errs { - if let Some(err) = err { - if err == &Error::Unexpected { - continue; - } - return Err(err.clone()); - } else { - all_at_eof = false; - } - } - - rustfs_io_metrics::record_stage_duration("sets_list_objects_list_merged", merge_started.elapsed().as_secs_f64() * 1000.0); - - debug!( - set_count = self.disk_set.len(), - all_at_eof = all_at_eof, - error_count = errs.iter().filter(|err| err.is_some()).count(), - "sets list_merged finished" - ); - - _ = all_at_eof; - Ok(Vec::new()) - } - - #[allow(unused_assignments)] - pub async fn walk_internal( - self: Arc, - rx: CancellationToken, - bucket: &str, - prefix: &str, - result: Sender, - opts: WalkOptions, - ) -> Result<()> { - check_list_objs_args(bucket, prefix, &None)?; - - let mut futures = Vec::new(); - let mut inputs = Vec::new(); - - for set in &self.disk_set { - let (mut disks, infos, _) = set.get_online_disks_with_healing_and_info(true).await; - let opts = opts.clone(); - let (sender, list_out_rx) = mpsc::channel::(1); - inputs.push(list_out_rx); - let rx_clone = rx.clone(); - let set = set.clone(); - futures.push(async move { - let mut ask_disks = get_list_quorum(&opts.ask_disks, set.set_drive_count as i32); - if ask_disks == -1 { - let new_disks = get_quorum_disks(&disks, &infos, disks.len().div_ceil(2)); - if !new_disks.is_empty() { - disks = new_disks; - } else { - ask_disks = get_list_quorum("strict", set.set_drive_count as i32); - } - } - - if set.set_drive_count == 4 { - ask_disks = bounded_usize_to_i32(disks.len()); - } else if ask_disks > bounded_usize_to_i32(disks.len()) { - ask_disks = clamp_ask_disks_to_available(ask_disks, disks.len()); - } - - let listing_quorum = listing_quorum_from_ask_disks(ask_disks); - let enforce_write_quorum = enforce_latest_listing_write_quorum(opts.latest_only, &opts.ask_disks); - let write_quorum_parity = set.default_parity_count; - let required_obj_quorum = latest_listing_required_object_quorum( - listing_quorum, - set.set_drive_count, - write_quorum_parity, - enforce_write_quorum, - ); - ask_disks = expand_ask_disks_for_object_quorum(ask_disks, disks.len(), required_obj_quorum); - let fallback_disks = if let Some(asked_disks) = positive_ask_disks(ask_disks) - && disks.len() > asked_disks - { - let mut rand = rand::rng(); - disks.shuffle(&mut rand); - disks.split_off(asked_disks) - } else { - Vec::new() - }; - let fallback_disks = Arc::new(fallback_disks); - let claim_tracker = FallbackClaimTracker::default(); - - let obj_quorum = - latest_listing_object_quorum(listing_quorum, set.set_drive_count, write_quorum_parity, enforce_write_quorum); - let raw_min_disks = latest_listing_raw_min_disks(listing_quorum, obj_quorum, enforce_write_quorum); - let resolver = list_metadata_resolution_params(bucket.to_owned(), listing_quorum, obj_quorum, !opts.latest_only); - let agreed_resolver = resolver.clone(); - let partial_resolver = resolver.clone(); - let reader_disks = disks.len(); - - let path = base_dir_from_prefix(prefix); - ensure_non_empty_listing_disks(bucket, &path, &disks)?; - - let mut filter_prefix = prefix - .trim_start_matches(&path) - .trim_start_matches(SLASH_SEPARATOR) - .trim_end_matches(SLASH_SEPARATOR) - .to_owned(); - if filter_prefix == path { - filter_prefix = "".to_owned(); - } - - let tx1 = sender.clone(); - let tx2 = sender.clone(); - let supplement = ListingSupplement::new( - ListingSupplementOptions { - bucket: bucket.to_owned(), - path: path.clone(), - recursive: true, - incl_deleted: !opts.latest_only, - filter_prefix: Some(filter_prefix.clone()), - forward_to: opts.marker.clone(), - per_disk_limit: bounded_usize_to_i32(opts.limit), - skip_total_timeout: opts.walkdir_timeout.is_none(), - walkdir_timeout: opts.walkdir_timeout, - walkdir_stall_timeout: opts.walkdir_stall_timeout, - }, - fallback_disks.clone(), - claim_tracker.clone(), - ); - let agreed_supplement = supplement.clone(); - let partial_supplement = supplement; - - list_path_raw_with_claim_tracker( - rx_clone, - ListPathRawOptions { - disks: disks.iter().cloned().map(Some).collect(), - fallback_disks: fallback_disks.iter().cloned().map(Some).collect(), - bucket: bucket.to_owned(), - path, - recursive: true, - incl_deleted: !opts.latest_only, - filter_prefix: Some(filter_prefix), - forward_to: opts.marker.clone(), - min_disks: raw_min_disks, - per_disk_limit: bounded_usize_to_i32(opts.limit), - // Skip the total walkdir timeout for listing operations. - // Large buckets (millions of objects) can take longer than - // the default 5s walkdir timeout to produce the first page - // of results. The stall timeout still protects against - // drives that stop making forward progress. - skip_walkdir_total_timeout: opts.walkdir_timeout.is_none(), - walkdir_timeout: opts.walkdir_timeout, - walkdir_stall_timeout: opts.walkdir_stall_timeout, - agreed: Some(Box::new(move |entry: MetaCacheEntry| { - Box::pin({ - let value = tx1.clone(); - let resolver = agreed_resolver.clone(); - let supplement = agreed_supplement.clone(); - async move { - let entry = match resolve_agreed_listing_entry( - entry, - reader_disks, - resolver.clone(), - enforce_write_quorum, - ) { - ListingEntryResolution::Resolved(entry) => entry, - ListingEntryResolution::NeedsSupplement(entry, fallback) => { - let Some(entry) = resolve_agreed_listing_entry_with_supplement( - entry, - reader_disks, - resolver, - enforce_write_quorum, - supplement, - ) - .await - .or(fallback) else { - return; - }; - entry - } - ListingEntryResolution::Rejected => return, - }; - if entry.is_dir() { - return; - } - if let Err(err) = value.send(entry).await { - error!("list_path send fail {:?}", err); - } - } - }) - })), - partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option]| { - Box::pin({ - let value = tx2.clone(); - let resolver = partial_resolver.clone(); - let supplement = partial_supplement.clone(); - async move { - if let Some(entry) = resolve_listing_entries_with_supplement( - entries, - resolver, - enforce_write_quorum, - supplement, - ) - .await - && let Err(err) = value.send(entry).await - { - error!("list_path send fail {:?}", err); - } - } - }) - })), - finished: None, - ..Default::default() - }, - claim_tracker, - ) - .await - }); - } - - let (merge_tx, mut merge_rx) = mpsc::channel::(100); - let bucket = bucket.to_owned(); - let bucket_clone = bucket.clone(); - - let vcf = match get_versioning_config(&bucket).await { - Ok((res, _)) => Some(res), - Err(_) => None, - }; - - tokio::spawn( - async move { - let mut sent_err = false; - while let Some(entry) = merge_rx.recv().await { - if opts.latest_only { - let fi = match entry.to_fileinfo(&bucket_clone) { - Ok(res) => res, - Err(err) => { - if !sent_err { - let item = ObjectInfoOrErr { - item: None, - err: Some(err.into()), - }; - if let Err(err) = result.send(item).await { - error!("walk result send err {:?}", err); - } - sent_err = true; - return; - } - continue; - } - }; - - if let Some(filter) = opts.filter { - if filter(&fi) { - let item = ObjectInfoOrErr { - item: Some(ObjectInfo::from_file_info(&fi, &bucket_clone, &fi.name, { - if let Some(v) = &vcf { v.versioned(&fi.name) } else { false } - })), - err: None, - }; - if let Err(err) = result.send(item).await { - error!("walk result send err {:?}", err); - } - } - } else { - let item = ObjectInfoOrErr { - item: Some(ObjectInfo::from_file_info(&fi, &bucket_clone, &fi.name, { - if let Some(v) = &vcf { v.versioned(&fi.name) } else { false } - })), - err: None, - }; - if let Err(err) = result.send(item).await { - error!("walk result send err {:?}", err); - } - } - continue; - } - - let fvs = match if opts.include_free_versions { - entry.file_info_versions_with_free_versions(&bucket_clone) - } else { - entry.file_info_versions(&bucket_clone) - } { - Ok(res) => res, - Err(err) => { - let item = ObjectInfoOrErr { - item: None, - err: Some(err.into()), - }; - if let Err(err) = result.send(item).await { - error!("walk result send err {:?}", err); - } - return; - } - }; - - for fi in &fvs.versions { - if let Some(filter) = opts.filter { - if filter(fi) { - let item = ObjectInfoOrErr { - item: Some(ObjectInfo::from_file_info(fi, &bucket_clone, &fi.name, { - if let Some(v) = &vcf { v.versioned(&fi.name) } else { false } - })), - err: None, - }; - if let Err(err) = result.send(item).await { - error!("walk result send err {:?}", err); - } - } - } else { - let item = ObjectInfoOrErr { - item: Some(ObjectInfo::from_file_info(fi, &bucket_clone, &fi.name, { - if let Some(v) = &vcf { v.versioned(&fi.name) } else { false } - })), - err: None, - }; - if let Err(err) = result.send(item).await { - error!("walk result send err {:?}", err); - } - } - } - - if opts.include_free_versions { - for fi in &fvs.free_versions { - if let Some(filter) = opts.filter { - if filter(fi) { - let item = ObjectInfoOrErr { - item: Some(ObjectInfo::from_file_info(fi, &bucket_clone, &fi.name, { - if let Some(v) = &vcf { v.versioned(&fi.name) } else { false } - })), - err: None, - }; - if let Err(err) = result.send(item).await { - error!("walk result send err {:?}", err); - } - } - } else { - let item = ObjectInfoOrErr { - item: Some(ObjectInfo::from_file_info(fi, &bucket_clone, &fi.name, { - if let Some(v) = &vcf { v.versioned(&fi.name) } else { false } - })), - err: None, - }; - if let Err(err) = result.send(item).await { - error!("walk result send err {:?}", err); - } - } - } - } - } - } - .instrument(tracing::Span::current()), - ); - - tokio::spawn( - async move { - if let Err(err) = merge_entry_channels(rx, inputs, merge_tx, 1).await { - error!("merge_entry_channels err {:?}", err) - } - } - .instrument(tracing::Span::current()), - ); - - let walk_started = std::time::Instant::now(); - let walk_results = join_all(futures).await; - let mut errs = Vec::new(); - for walk_result in walk_results { - match walk_result { - Ok(()) => errs.push(None), - Err(err) => errs.push(Some(err.into())), - } - } - rustfs_io_metrics::record_stage_duration( - "sets_list_objects_walk_internal", - walk_started.elapsed().as_secs_f64() * 1000.0, - ); - - let result = walk_result_from_set_errors(&errs); - if let Err(err) = &result { - error!( - bucket = %bucket, - prefix = %prefix, - error = ?err, - set_errors = ?errs, - "walk_internal list_path_raw tasks failed" - ); - } - - result - } -} - impl SetDisks { pub(crate) async fn inner_list_object_versions_for_recursive_delete( self: Arc, diff --git a/scripts/check_logging_guardrails.sh b/scripts/check_logging_guardrails.sh index a5095d8b3..c9b0a14eb 100755 --- a/scripts/check_logging_guardrails.sh +++ b/scripts/check_logging_guardrails.sh @@ -985,7 +985,9 @@ trace_hot_spans=( "crates/ecstore/src/set_disk/ops/object.rs:get_object_info" "crates/ecstore/src/store/mod.rs:list_objects_v2" "crates/ecstore/src/store/list.rs:handle_list_objects_v2" - "crates/ecstore/src/core/sets.rs:list_objects_v2" + # The pool-level Sets::list_objects_v2 wrapper was removed with its duplicate + # pagination pipeline (backlog#1821); the remaining ECStore and SetDisks + # wrappers below still carry the TRACE requirement. "crates/ecstore/src/set_disk/ops/list.rs:list_objects_v2" "rustfs/src/app/bucket_usecase.rs:execute_list_objects_v2" "rustfs/src/app/object_usecase.rs:execute_get_object"