fix(ecstore): preserve Sets listing compatibility

This commit is contained in:
overtrue
2026-08-14 11:02:01 +08:00
parent 5aa72c30b2
commit 06c6a2a914
4 changed files with 911 additions and 11 deletions
+71 -8
View File
@@ -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<ObjectInfo>;
type ListObjectVersionsInfo = StorageListObjectVersionsInfo<ObjectInfo>;
type ObjectInfoOrErr = StorageObjectInfoOrErr<ObjectInfo, Error>;
type WalkOptions = StorageWalkOptions<fn(&FileInfo) -> 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<ObjectInfoOrErr>;
#[tracing::instrument(level = "trace", skip(self))]
async fn list_objects_v2(
self: Arc<Self>,
bucket: &str,
prefix: &str,
continuation_token: Option<String>,
delimiter: Option<String>,
max_keys: i32,
fetch_owner: bool,
start_after: Option<String>,
incl_deleted: bool,
) -> Result<ListObjectsV2Info> {
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<Self>,
bucket: &str,
prefix: &str,
marker: Option<String>,
version_marker: Option<String>,
delimiter: Option<String>,
max_keys: i32,
) -> Result<ListObjectVersionsInfo> {
self.inner_list_object_versions(bucket, prefix, marker, version_marker, delimiter, max_keys)
.await
}
async fn walk(
self: Arc<Self>,
rx: CancellationToken,
bucket: &str,
prefix: &str,
result: tokio::sync::mpsc::Sender<ObjectInfoOrErr>,
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);
+827
View File
@@ -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<Self>,
bucket: &str,
prefix: &str,
continuation_token: Option<String>,
delimiter: Option<String>,
max_keys: i32,
_fetch_owner: bool,
start_after: Option<String>,
incl_deleted: bool,
) -> Result<ListObjectsV2Info> {
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<Self>,
bucket: &str,
prefix: &str,
marker: Option<String>,
delimiter: Option<String>,
max_keys: i32,
incl_deleted: bool,
) -> Result<ListObjectsInfo> {
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<Self>,
bucket: &str,
prefix: &str,
marker: Option<String>,
version_marker: Option<String>,
delimiter: Option<String>,
max_keys: i32,
) -> Result<ListObjectVersionsInfo> {
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<Self>, o: &ListPathOptions) -> Result<MetaCacheEntriesSortedResult> {
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::<Arc<Error>>(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<MetaCacheEntry>,
) -> Result<Vec<ObjectInfo>> {
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<Self>,
rx: CancellationToken,
bucket: &str,
prefix: &str,
result: Sender<ObjectInfoOrErr>,
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::<MetaCacheEntry>(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<DiskError>]| {
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::<MetaCacheEntry>(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<Self>,
@@ -157,6 +157,18 @@ fn ecstore_implements_storage_list_operations_contract() {
assert!(storage_list_operations_type_name::<ECStore>().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::<ECStore>().ends_with("::ECStore"));
+1 -3
View File
@@ -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"