mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-15 01:23:12 +00:00
chore(ecstore): remove the pool-level ListObjects pagination copy
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).
This commit is contained in:
@@ -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<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;
|
||||
@@ -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);
|
||||
|
||||
@@ -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<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>,
|
||||
|
||||
@@ -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"
|
||||
|
||||
Reference in New Issue
Block a user