From 06c6a2a914134d0c24bb1b194411c776bdf3bab0 Mon Sep 17 00:00:00 2001 From: overtrue Date: Fri, 14 Aug 2026 11:02:01 +0800 Subject: [PATCH] fix(ecstore): preserve Sets listing compatibility --- crates/ecstore/src/core/sets.rs | 79 +- crates/ecstore/src/store/list_objects.rs | 827 ++++++++++++++++++ .../tests/ecstore_contract_compat_test.rs | 12 + scripts/check_logging_guardrails.sh | 4 +- 4 files changed, 911 insertions(+), 11 deletions(-) diff --git a/crates/ecstore/src/core/sets.rs b/crates/ecstore/src/core/sets.rs index 7ef0f65b9..acb8b53b9 100644 --- a/crates/ecstore/src/core/sets.rs +++ b/crates/ecstore/src/core/sets.rs @@ -19,6 +19,7 @@ use crate::layout::set_heal::{formats_to_drives_info, new_heal_format_sets}; use crate::multipart_listing::paginate_multipart_listing; use crate::storage_api_contracts::{ bucket::{BucketInfo, BucketOperations, BucketOptions, DeleteBucketOptions, MakeBucketOptions}, + list::{StorageListObjectVersionsInfo, StorageListObjectsV2Info, StorageObjectInfoOrErr, StorageWalkOptions}, multipart::{CompletePart, ListMultipartsInfo, ListPartsInfo, MultipartInfo, MultipartUploadResult, PartInfo}, object::{DeletedObject, ObjectIO as _, ObjectOperations as _, ObjectToDelete}, range::HTTPRangeSpec, @@ -58,10 +59,16 @@ 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; +type ListObjectsV2Info = StorageListObjectsV2Info; +type ListObjectVersionsInfo = StorageListObjectVersionsInfo; +type ObjectInfoOrErr = StorageObjectInfoOrErr; +type WalkOptions = StorageWalkOptions bool>; + const LIST_MULTIPART_SETS_CONCURRENCY: usize = 4; fn is_idempotent_delete_prefix_error(err: &Error) -> bool { @@ -749,6 +756,67 @@ 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; @@ -1835,7 +1903,7 @@ mod tests { #[tokio::test(flavor = "multi_thread")] #[serial] - async fn set_level_list_objects_v2_lists_objects_written_through_the_pool() { + async fn sets_list_objects_v2_lists_objects_within_the_pool() { let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::Erasure).await; let format = FormatV3::new(1, 2); let mut endpoints = Vec::new(); @@ -1917,16 +1985,11 @@ mod tests { .await .expect("object should be written"); - // 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] + let result = sets .clone() .list_objects_v2(&bucket, "", None, None, 1000, false, None, false) .await - .expect("set-level listing should succeed"); + .expect("pool-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 9adadbb13..4ea8fc78a 100644 --- a/crates/ecstore/src/store/list_objects.rs +++ b/crates/ecstore/src/store/list_objects.rs @@ -16,6 +16,7 @@ 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::{ @@ -5029,6 +5030,832 @@ 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/crates/ecstore/tests/ecstore_contract_compat_test.rs b/crates/ecstore/tests/ecstore_contract_compat_test.rs index 496eab291..d1d303f7a 100644 --- a/crates/ecstore/tests/ecstore_contract_compat_test.rs +++ b/crates/ecstore/tests/ecstore_contract_compat_test.rs @@ -157,6 +157,18 @@ fn ecstore_implements_storage_list_operations_contract() { assert!(storage_list_operations_type_name::().ends_with("::ECStore")); } +#[test] +fn ecstore_pools_expose_storage_list_operations_contract() { + fn assert_contract(store: &ECStore) { + let future = store.pools[0] + .clone() + .list_objects_v2("bucket", "", None, None, 1, false, None, false); + drop(future); + } + + let _ = assert_contract; +} + #[test] fn ecstore_implements_storage_multipart_operations_contract() { assert!(storage_multipart_operations_type_name::().ends_with("::ECStore")); diff --git a/scripts/check_logging_guardrails.sh b/scripts/check_logging_guardrails.sh index 6db62ecc2..eb0640635 100755 --- a/scripts/check_logging_guardrails.sh +++ b/scripts/check_logging_guardrails.sh @@ -987,9 +987,7 @@ trace_hot_spans=( # The ECStore handle_list_objects_v2 forwarder was folded into the trait impl # above, so store/mod.rs now carries this hot path's TRACE requirement # directly (backlog#1821). - # 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/core/sets.rs:list_objects_v2" "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"