mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 12:49:04 +00:00
Merge remote-tracking branch 'origin/main' into cxymds/perf-decommission-entry-budget
# Conflicts: # crates/ecstore/src/core/pools.rs
This commit is contained in:
@@ -49,4 +49,3 @@ Applies to `crates/ecstore/`.
|
||||
## Suggested Validation
|
||||
|
||||
- `cargo test -p rustfs-ecstore`
|
||||
- Full gate before commit: `make pre-commit`
|
||||
|
||||
@@ -866,7 +866,7 @@ impl BucketTargetSys {
|
||||
return Some(cli);
|
||||
}
|
||||
|
||||
// TODO: spawn a task to reload the target
|
||||
// TODO(backlog): spawn an async task to proactively reload the replication target
|
||||
if self.is_reloading_target(bucket, arn).await {
|
||||
return None;
|
||||
}
|
||||
|
||||
@@ -454,7 +454,7 @@ impl S3PeerSys {
|
||||
}
|
||||
}
|
||||
topology_complete &= bucket_map.values().all(|count| *count >= quorum);
|
||||
// TODO: MRF
|
||||
// TODO(backlog): integrate MRF backlog stats into scanner bucket listing
|
||||
}
|
||||
|
||||
let mut buckets: Vec<BucketInfo> = result_map.into_values().collect();
|
||||
|
||||
@@ -2406,7 +2406,7 @@ impl DiskAPI for RemoteDisk {
|
||||
return errors;
|
||||
}
|
||||
|
||||
// TODO: use Error not string
|
||||
// TODO(backlog): replace string errors with typed `StorageError` variants
|
||||
|
||||
let result = self
|
||||
.execute_with_timeout(
|
||||
|
||||
@@ -36,7 +36,8 @@ use crate::disk::error::DiskError;
|
||||
use crate::disk::{BUCKET_META_PREFIX, RUSTFS_META_BUCKET};
|
||||
use crate::error::{Error, Result};
|
||||
use crate::error::{
|
||||
StorageError, is_err_bucket_exists, is_err_bucket_not_found, is_err_object_not_found, is_err_version_not_found,
|
||||
StorageError, is_err_bucket_exists, is_err_bucket_not_found, is_err_object_not_found, is_err_operation_canceled,
|
||||
is_err_version_not_found,
|
||||
};
|
||||
use crate::layout::endpoints::EndpointServerPools;
|
||||
use crate::object_api::{GetObjectReader, ObjectOptions};
|
||||
@@ -817,7 +818,60 @@ async fn load_decommission_entry_exact_versions(
|
||||
}
|
||||
|
||||
fn resolve_decommission_check_after_list_result(list_result: Result<()>, entry_error: Option<Error>) -> Result<()> {
|
||||
if let Some(err) = entry_error { Err(err) } else { list_result }
|
||||
match list_result {
|
||||
Ok(()) => entry_error.map_or(Ok(()), Err),
|
||||
Err(list_err) => resolve_decommission_listing_error(Some(list_err), entry_error).map_or(Ok(()), Err),
|
||||
}
|
||||
}
|
||||
|
||||
fn resolve_decommission_listing_error(listing_error: Option<Error>, entry_error: Option<Error>) -> Option<Error> {
|
||||
match (listing_error, entry_error) {
|
||||
(Some(listing_error), Some(entry_error)) if is_err_operation_canceled(&listing_error) => Some(entry_error),
|
||||
(Some(listing_error), Some(entry_error)) if is_err_operation_canceled(&entry_error) => Some(listing_error),
|
||||
(Some(listing_error), _) => Some(listing_error),
|
||||
(None, entry_error) => entry_error,
|
||||
}
|
||||
}
|
||||
|
||||
fn decommission_unresolved_listing_error(
|
||||
bucket: &str,
|
||||
prefix: &str,
|
||||
candidate: Option<&str>,
|
||||
candidate_count: usize,
|
||||
disk_error_count: usize,
|
||||
pool_index: usize,
|
||||
set_index: usize,
|
||||
) -> Error {
|
||||
let location = candidate.unwrap_or(prefix);
|
||||
Error::other(format!(
|
||||
"decommission listing could not resolve metadata for {bucket}/{location} on pool {pool_index} set {set_index} ({candidate_count} candidate(s), {disk_error_count} disk error(s))"
|
||||
))
|
||||
}
|
||||
|
||||
fn resolve_decommission_partial_listing_entry(
|
||||
entries: MetaCacheEntries,
|
||||
resolver: MetadataResolutionParams,
|
||||
bucket: &str,
|
||||
prefix: &str,
|
||||
disk_error_count: usize,
|
||||
pool_index: usize,
|
||||
set_index: usize,
|
||||
) -> Result<MetaCacheEntry> {
|
||||
let candidate_count = entries.as_ref().iter().flatten().count();
|
||||
if let Some(entry) = entries.resolve(resolver) {
|
||||
return Ok(entry);
|
||||
}
|
||||
|
||||
let candidate = entries.as_ref().iter().flatten().map(|entry| entry.name.as_str()).next();
|
||||
Err(decommission_unresolved_listing_error(
|
||||
bucket,
|
||||
prefix,
|
||||
candidate,
|
||||
candidate_count,
|
||||
disk_error_count,
|
||||
pool_index,
|
||||
set_index,
|
||||
))
|
||||
}
|
||||
|
||||
fn resolve_decommission_pool_meta_reload_result(result: Result<()>, stage: &str) -> Result<()> {
|
||||
@@ -2976,6 +3030,34 @@ impl ECStore {
|
||||
canceled_worker
|
||||
}
|
||||
|
||||
async fn clear_decommission_canceler_for_generation(&self, idx: usize, generation: OffsetDateTime) {
|
||||
let _start_guard = self.start_gate.lock().await;
|
||||
let generation_is_current = self
|
||||
.pool_meta
|
||||
.read()
|
||||
.await
|
||||
.pools
|
||||
.get(idx)
|
||||
.and_then(|pool| pool.decommission.as_ref())
|
||||
.and_then(|info| info.start_time)
|
||||
== Some(generation);
|
||||
if !generation_is_current {
|
||||
return;
|
||||
}
|
||||
|
||||
let mut cancelers = self.decommission_cancelers.write().await;
|
||||
if take_decommission_canceler(cancelers.as_mut_slice(), idx).is_none() {
|
||||
warn!(
|
||||
event = EVENT_DECOMMISSION_STATE,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_POOLS,
|
||||
pool_index = idx,
|
||||
state = "canceler_already_cleared",
|
||||
"Decommission canceler already cleared"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
async fn cancel_decommission_routines_and_wait(&self, indices: &[usize]) {
|
||||
let canceled_worker = {
|
||||
let mut cancelers = self.decommission_cancelers.write().await;
|
||||
@@ -3340,6 +3422,7 @@ impl ECStore {
|
||||
let list_rx_for_drain = list_rx.clone();
|
||||
let list_bi = bi.clone();
|
||||
let list_outstanding = outstanding.clone();
|
||||
let list_entry_error = entry_error.clone();
|
||||
let mut listing = tokio::spawn(async move {
|
||||
run_decommission_listing_with_retry_and_drain(
|
||||
list_rx.clone(),
|
||||
@@ -3352,7 +3435,11 @@ impl ECStore {
|
||||
let set = list_set.clone();
|
||||
let rx = list_rx_for_list.clone();
|
||||
let bucket = list_bi.clone();
|
||||
async move { set.list_objects_to_decommission(rx, bucket, callback).await }
|
||||
let entry_error = list_entry_error.clone();
|
||||
async move {
|
||||
set.list_objects_to_decommission(rx, bucket, callback, entry_error, idx, set_idx)
|
||||
.await
|
||||
}
|
||||
},
|
||||
move || {
|
||||
let rx = list_rx_for_drain.clone();
|
||||
@@ -4071,20 +4158,6 @@ impl ECStore {
|
||||
idx: usize,
|
||||
entry_budget: Arc<Semaphore>,
|
||||
) -> Result<()> {
|
||||
defer!(|| async {
|
||||
let mut cancelers = self.decommission_cancelers.write().await;
|
||||
if take_decommission_canceler(cancelers.as_mut_slice(), idx).is_none() {
|
||||
warn!(
|
||||
event = EVENT_DECOMMISSION_STATE,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_POOLS,
|
||||
pool_index = idx,
|
||||
state = "canceler_already_cleared",
|
||||
"Decommission canceler already cleared"
|
||||
);
|
||||
}
|
||||
});
|
||||
|
||||
if let Err(err) = self.promote_queued_decommission(idx).await {
|
||||
resolve_decommission_terminal_mark_after_error_result(self.decommission_failed(idx).await, idx, &err)?;
|
||||
return Err(err);
|
||||
@@ -4112,6 +4185,9 @@ impl ECStore {
|
||||
return Ok(());
|
||||
}
|
||||
let generation = self.active_decommission_generation(idx).await?;
|
||||
defer!(|| async {
|
||||
self.clear_decommission_canceler_for_generation(idx, generation).await;
|
||||
});
|
||||
let result = self.decommission_in_background(rx.clone(), idx, entry_budget).await;
|
||||
|
||||
let (final_state, canceled, cmd_line) = {
|
||||
@@ -4166,7 +4242,11 @@ impl ECStore {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
resolve_decommission_terminal_mark_after_error_result(self.decommission_failed(idx).await, idx, &err)?;
|
||||
resolve_decommission_terminal_mark_after_error_result(
|
||||
self.decommission_failed_with_generation(idx, Some(generation)).await,
|
||||
idx,
|
||||
&err,
|
||||
)?;
|
||||
warn!(
|
||||
event = EVENT_DECOMMISSION_STATE,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
@@ -4217,7 +4297,11 @@ impl ECStore {
|
||||
rx.cancel();
|
||||
return Ok(());
|
||||
}
|
||||
resolve_decommission_terminal_mark_result(self.decommission_failed(idx).await, "failed", &cmd_line)?;
|
||||
resolve_decommission_terminal_mark_result(
|
||||
self.decommission_failed_with_generation(idx, Some(generation)).await,
|
||||
"failed",
|
||||
&cmd_line,
|
||||
)?;
|
||||
return Err(Error::other(format!(
|
||||
"failed to finalize decommission for pool {cmd_line}: post-check failed: {err}"
|
||||
)));
|
||||
@@ -4238,7 +4322,11 @@ impl ECStore {
|
||||
state = "marking_completed",
|
||||
"Decommission marking completed state"
|
||||
);
|
||||
resolve_decommission_terminal_mark_result(self.complete_decommission(idx).await, "completed", &cmd_line)?;
|
||||
resolve_decommission_terminal_mark_result(
|
||||
self.complete_decommission_with_generation(idx, Some(generation)).await,
|
||||
"completed",
|
||||
&cmd_line,
|
||||
)?;
|
||||
}
|
||||
DecommissionFinalState::Failed => {
|
||||
warn!(
|
||||
@@ -4250,7 +4338,11 @@ impl ECStore {
|
||||
state = "marking_failed",
|
||||
"Decommission marking failed state"
|
||||
);
|
||||
resolve_decommission_terminal_mark_result(self.decommission_failed(idx).await, "failed", &cmd_line)?;
|
||||
resolve_decommission_terminal_mark_result(
|
||||
self.decommission_failed_with_generation(idx, Some(generation)).await,
|
||||
"failed",
|
||||
&cmd_line,
|
||||
)?;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -4268,16 +4360,28 @@ impl ECStore {
|
||||
|
||||
#[tracing::instrument(skip(self))]
|
||||
pub async fn decommission_failed(&self, idx: usize) -> Result<()> {
|
||||
self.decommission_failed_with_generation(idx, None).await
|
||||
}
|
||||
|
||||
async fn decommission_failed_with_generation(&self, idx: usize, generation: Option<OffsetDateTime>) -> Result<()> {
|
||||
ensure_decommission_terminal_operation_supported(self.single_pool(), "mark decommission failed")?;
|
||||
let _start_guard = self.start_gate.lock().await;
|
||||
|
||||
let (should_reload_pool_meta, previous_pool_meta) = {
|
||||
let mut pool_meta = self.pool_meta.write().await;
|
||||
if let Some(generation) = generation {
|
||||
match ensure_decommission_generation(&pool_meta, idx, generation) {
|
||||
Ok(()) => {}
|
||||
Err(Error::OperationCanceled) => return Ok(()),
|
||||
Err(err) => return Err(err),
|
||||
}
|
||||
}
|
||||
let previous_pool_meta = pool_meta.clone();
|
||||
let changed = pool_meta.decommission_failed(idx);
|
||||
(changed, changed.then_some(previous_pool_meta))
|
||||
};
|
||||
|
||||
{
|
||||
if should_reload_pool_meta {
|
||||
let mut cancelers = self.decommission_cancelers.write().await;
|
||||
take_and_cancel_decommission_canceler(cancelers.as_mut_slice(), idx);
|
||||
}
|
||||
@@ -4333,18 +4437,29 @@ impl ECStore {
|
||||
|
||||
#[tracing::instrument(skip(self))]
|
||||
pub async fn complete_decommission(&self, idx: usize) -> Result<()> {
|
||||
self.complete_decommission_with_generation(idx, None).await
|
||||
}
|
||||
|
||||
async fn complete_decommission_with_generation(&self, idx: usize, generation: Option<OffsetDateTime>) -> Result<()> {
|
||||
ensure_decommission_terminal_operation_supported(self.single_pool(), "complete decommission")?;
|
||||
|
||||
let _start_guard = self.start_gate.lock().await;
|
||||
|
||||
let (should_reload_pool_meta, previous_pool_meta) = {
|
||||
let mut pool_meta = self.pool_meta.write().await;
|
||||
if let Some(generation) = generation {
|
||||
match ensure_decommission_generation(&pool_meta, idx, generation) {
|
||||
Ok(()) => {}
|
||||
Err(Error::OperationCanceled) => return Ok(()),
|
||||
Err(err) => return Err(err),
|
||||
}
|
||||
}
|
||||
let previous_pool_meta = pool_meta.clone();
|
||||
let changed = pool_meta.decommission_complete(idx);
|
||||
(changed, changed.then_some(previous_pool_meta))
|
||||
};
|
||||
|
||||
{
|
||||
if should_reload_pool_meta {
|
||||
let mut cancelers = self.decommission_cancelers.write().await;
|
||||
take_and_cancel_decommission_canceler(cancelers.as_mut_slice(), idx);
|
||||
}
|
||||
@@ -4676,7 +4791,7 @@ impl ECStore {
|
||||
let buckets = self.get_buckets_to_decommission().await?;
|
||||
let pool = self.pools[idx].clone();
|
||||
|
||||
for set in &pool.disk_set {
|
||||
for (set_index, set) in pool.disk_set.iter().enumerate() {
|
||||
for bucket_info in &buckets {
|
||||
let mut lifecycle_config = None;
|
||||
let mut object_lock_config = None;
|
||||
@@ -4771,7 +4886,7 @@ impl ECStore {
|
||||
});
|
||||
|
||||
let list_result = set
|
||||
.list_objects_to_decommission(callback_rx, bucket_info.clone(), callback)
|
||||
.list_objects_to_decommission(callback_rx, bucket_info.clone(), callback, entry_error.clone(), idx, set_index)
|
||||
.await;
|
||||
let entry_error = entry_error.lock().await.clone();
|
||||
resolve_decommission_check_after_list_result(list_result, entry_error)?;
|
||||
@@ -5575,20 +5690,27 @@ async fn record_decommission_entry_error(
|
||||
rx: &CancellationToken,
|
||||
err: Error,
|
||||
) {
|
||||
if rx.is_cancelled() {
|
||||
return;
|
||||
}
|
||||
|
||||
let mut first_err = entry_error.lock().await;
|
||||
if first_err.is_none() {
|
||||
if first_err.is_none() && !rx.is_cancelled() {
|
||||
*first_err = Some(err);
|
||||
rx.cancel();
|
||||
}
|
||||
}
|
||||
|
||||
impl SetDisks {
|
||||
#[tracing::instrument(skip(self, rx, cb_func))]
|
||||
#[tracing::instrument(skip(self, rx, cb_func, entry_error))]
|
||||
async fn list_objects_to_decommission(
|
||||
self: &Arc<Self>,
|
||||
rx: CancellationToken,
|
||||
bucket_info: DecomBucketInfo,
|
||||
cb_func: ListCallback,
|
||||
entry_error: Arc<tokio::sync::Mutex<Option<Error>>>,
|
||||
pool_index: usize,
|
||||
set_index: usize,
|
||||
) -> Result<()> {
|
||||
let (disks, _) = self.get_online_disks_with_healing(false).await;
|
||||
ensure_decommission_listing_disks_available(!disks.is_empty(), &bucket_info.name)?;
|
||||
@@ -5603,6 +5725,12 @@ impl SetDisks {
|
||||
};
|
||||
|
||||
let cb1 = cb_func.clone();
|
||||
let unresolved_error = entry_error.clone();
|
||||
let unresolved_rx = rx.clone();
|
||||
let unresolved_bucket = bucket_info.name.clone();
|
||||
let unresolved_prefix = bucket_info.prefix.clone();
|
||||
let unresolved_pool_index = pool_index;
|
||||
let unresolved_set_index = set_index;
|
||||
|
||||
list_path_raw(
|
||||
rx,
|
||||
@@ -5615,20 +5743,51 @@ impl SetDisks {
|
||||
skip_walkdir_total_timeout: true,
|
||||
walkdir_stall_timeout: Some(DECOMMISSION_BACKGROUND_WALKDIR_STALL_TIMEOUT),
|
||||
agreed: Some(Box::new(move |entry: MetaCacheEntry| Box::pin(cb1(entry)))),
|
||||
partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option<DiskError>]| {
|
||||
partial: Some(Box::new(move |entries: MetaCacheEntries, errs: &[Option<DiskError>]| {
|
||||
let resolver = resolver.clone();
|
||||
let cb_func = cb_func.clone();
|
||||
match entries.resolve(resolver) {
|
||||
Some(entry) => {
|
||||
let bucket = unresolved_bucket.clone();
|
||||
let prefix = unresolved_prefix.clone();
|
||||
let unresolved_error = unresolved_error.clone();
|
||||
let unresolved_rx = unresolved_rx.clone();
|
||||
let pool_index = unresolved_pool_index;
|
||||
let set_index = unresolved_set_index;
|
||||
let disk_error_count = errs.iter().flatten().count();
|
||||
if unresolved_rx.is_cancelled() {
|
||||
return Box::pin(async {});
|
||||
}
|
||||
|
||||
match resolve_decommission_partial_listing_entry(
|
||||
entries,
|
||||
resolver,
|
||||
&bucket,
|
||||
&prefix,
|
||||
disk_error_count,
|
||||
pool_index,
|
||||
set_index,
|
||||
) {
|
||||
Ok(entry) => {
|
||||
warn!("decommission_pool: list_objects_to_decommission get {}", &entry.name);
|
||||
Box::pin(async move {
|
||||
cb_func(entry).await;
|
||||
})
|
||||
}
|
||||
None => {
|
||||
warn!("decommission_pool: list_objects_to_decommission get none");
|
||||
Box::pin(async {})
|
||||
}
|
||||
Err(err) => Box::pin(async move {
|
||||
if unresolved_rx.is_cancelled() {
|
||||
return;
|
||||
}
|
||||
warn!(
|
||||
event = EVENT_DECOMMISSION_BUCKET,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_POOLS,
|
||||
bucket = %bucket,
|
||||
prefix = %prefix,
|
||||
state = "unresolved_entry",
|
||||
error = %err,
|
||||
"Decommission listing failed closed on unresolved metadata"
|
||||
);
|
||||
record_decommission_entry_error(&unresolved_error, &unresolved_rx, err).await;
|
||||
}),
|
||||
}
|
||||
})),
|
||||
..Default::default()
|
||||
@@ -5636,6 +5795,10 @@ impl SetDisks {
|
||||
)
|
||||
.await?;
|
||||
|
||||
if let Some(err) = entry_error.lock().await.clone() {
|
||||
return Err(err);
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@@ -5844,11 +6007,12 @@ mod pools_tests {
|
||||
get_by_index, has_active_decommission_canceler, is_decommission_active, is_decommission_cancel_requested,
|
||||
load_decommission_entry_versions, local_decommission_queue_prefix, mark_decommission_bucket_done,
|
||||
merge_pool_status_refresh, missing_decommission_worker_prefix, observe_decommission_terminal_reload_result,
|
||||
pool_meta_has_active_decommission, require_decommission_store, resolve_decommission_bucket_done_save_result,
|
||||
resolve_decommission_bucket_state, resolve_decommission_check_after_list_result,
|
||||
resolve_decommission_entry_cleanup_delete_result, resolve_decommission_entry_exact_versions,
|
||||
resolve_decommission_entry_reload_result, resolve_decommission_listing_worker_result,
|
||||
resolve_decommission_optional_bucket_config_result, resolve_decommission_pool_meta_reload_result,
|
||||
pool_meta_has_active_decommission, record_decommission_entry_error, require_decommission_store,
|
||||
resolve_decommission_bucket_done_save_result, resolve_decommission_bucket_state,
|
||||
resolve_decommission_check_after_list_result, resolve_decommission_entry_cleanup_delete_result,
|
||||
resolve_decommission_entry_exact_versions, resolve_decommission_entry_reload_result, resolve_decommission_listing_error,
|
||||
resolve_decommission_listing_worker_result, resolve_decommission_optional_bucket_config_result,
|
||||
resolve_decommission_partial_listing_entry, resolve_decommission_pool_meta_reload_result,
|
||||
resolve_decommission_preflight_heal_result, resolve_decommission_progress_save_result,
|
||||
resolve_decommission_spawn_failure_result, resolve_decommission_terminal_mark_after_error_result,
|
||||
resolve_decommission_terminal_mark_result, resolve_decommission_update_after_result,
|
||||
@@ -5867,7 +6031,9 @@ mod pools_tests {
|
||||
use crate::error::{Error, StorageError};
|
||||
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
|
||||
use crate::services::rebalance::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats};
|
||||
use rustfs_filemeta::{FileInfo, FileInfoVersions, MetaCacheEntry, ObjectPartInfo};
|
||||
use rustfs_filemeta::{
|
||||
FileInfo, FileInfoVersions, MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams, ObjectPartInfo,
|
||||
};
|
||||
use rustfs_rio::Index;
|
||||
use std::sync::{
|
||||
Arc,
|
||||
@@ -7066,6 +7232,65 @@ mod pools_tests {
|
||||
assert!(matches!(err, Error::SlowDown));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_resolve_decommission_partial_listing_entry_rejects_unresolved_metadata() {
|
||||
let err = resolve_decommission_partial_listing_entry(
|
||||
MetaCacheEntries(vec![None]),
|
||||
MetadataResolutionParams {
|
||||
dir_quorum: 2,
|
||||
obj_quorum: 2,
|
||||
bucket: "bucket-a".to_string(),
|
||||
..Default::default()
|
||||
},
|
||||
"bucket-a",
|
||||
"prefix/",
|
||||
1,
|
||||
2,
|
||||
3,
|
||||
)
|
||||
.expect_err("unresolved partial listing must fail closed");
|
||||
|
||||
let message = err.to_string();
|
||||
assert!(message.contains("decommission listing could not resolve metadata"));
|
||||
assert!(message.contains("bucket-a/prefix/"));
|
||||
assert!(message.contains("pool 2 set 3"));
|
||||
assert!(message.contains("1 disk error(s)"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_record_decommission_entry_error_cancels_listing_and_preserves_first_error() {
|
||||
let entry_error = Arc::new(tokio::sync::Mutex::new(None));
|
||||
let rx = CancellationToken::new();
|
||||
|
||||
record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await;
|
||||
record_decommission_entry_error(&entry_error, &rx, Error::OperationCanceled).await;
|
||||
|
||||
assert!(rx.is_cancelled());
|
||||
assert!(matches!(*entry_error.lock().await, Some(Error::SlowDown)));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_record_decommission_entry_error_ignores_already_canceled_listing() {
|
||||
let entry_error = Arc::new(tokio::sync::Mutex::new(None));
|
||||
let rx = CancellationToken::new();
|
||||
rx.cancel();
|
||||
|
||||
record_decommission_entry_error(&entry_error, &rx, Error::SlowDown).await;
|
||||
|
||||
assert!(entry_error.lock().await.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_resolve_decommission_listing_error_preserves_real_listing_failure() {
|
||||
let err = resolve_decommission_listing_error(Some(Error::SlowDown), Some(Error::OperationCanceled))
|
||||
.expect("listing failure should be returned");
|
||||
assert!(matches!(err, Error::SlowDown));
|
||||
|
||||
let err = resolve_decommission_listing_error(Some(Error::OperationCanceled), Some(Error::SlowDown))
|
||||
.expect("entry failure should be returned");
|
||||
assert!(matches!(err, Error::SlowDown));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_resolve_decommission_check_after_list_result_returns_list_result_without_entry_error() {
|
||||
let err = resolve_decommission_check_after_list_result(Err(Error::OperationCanceled), None)
|
||||
|
||||
@@ -249,7 +249,7 @@ impl Sets {
|
||||
|
||||
self.connect_disks().await;
|
||||
|
||||
// TODO: config interval
|
||||
// TODO(backlog): make monitor_and_connect interval configurable instead of hardcoded 15s
|
||||
let mut interval = tokio::time::interval(Duration::from_secs(15));
|
||||
loop {
|
||||
tokio::select! {
|
||||
|
||||
@@ -5215,8 +5215,8 @@ impl LocalDisk {
|
||||
|
||||
let cache = Cache::new(update_fn, Duration::from_secs(1), Opts::default());
|
||||
|
||||
// TODO: DIRECT support
|
||||
// TODD: DiskInfo
|
||||
// TODO(backlog): add O_DIRECT I/O support for performance-critical paths
|
||||
// TODO(backlog): populate DiskInfo in constructor
|
||||
let mut disk = Self {
|
||||
root: root.clone(),
|
||||
publication_root,
|
||||
@@ -5751,7 +5751,7 @@ impl LocalDisk {
|
||||
|
||||
// return Ok(());
|
||||
|
||||
// TODO: async notifications for disk space checks and trash cleanup
|
||||
// TODO(backlog): make disk space checks and trash cleanup event-driven instead of poll-based
|
||||
|
||||
let trash_path = self.io_get_object_path(RUSTFS_META_TMP_DELETED_BUCKET, Uuid::new_v4().to_string().as_str())?;
|
||||
// if let Some(parent) = trash_path.parent() {
|
||||
@@ -5997,7 +5997,7 @@ impl LocalDisk {
|
||||
|
||||
#[hotpath::measure(impl_type = "LocalDisk")]
|
||||
async fn read_all_data(&self, volume: &str, volume_dir: impl AsRef<Path>, file_path: impl AsRef<Path>) -> Result<Vec<u8>> {
|
||||
// TODO: timeout support
|
||||
// TODO(backlog): add configurable timeout for read_all_data operations
|
||||
let (data, _) = self.read_all_data_with_dmtime(volume, volume_dir, file_path).await?;
|
||||
Ok(data)
|
||||
}
|
||||
@@ -6674,7 +6674,7 @@ impl LocalDisk {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
// TODO: add lock
|
||||
// TODO(backlog): add directory listing lock to prevent concurrent enumeration
|
||||
|
||||
let stall = opts.stall_timeout_duration();
|
||||
|
||||
@@ -8796,7 +8796,7 @@ impl DiskAPI for LocalDisk {
|
||||
Ok(entries)
|
||||
}
|
||||
|
||||
// FIXME: TODO: io.writer TODO cancel
|
||||
// TODO(backlog): support io.writer cancellation and early termination in walk_dir
|
||||
#[tracing::instrument(level = "trace", skip_all)]
|
||||
async fn walk_dir<W: AsyncWrite + Unpin + Send>(&self, opts: WalkDirOptions, wr: &mut W) -> Result<()> {
|
||||
self.wait_for_startup_cleanup().await;
|
||||
@@ -9880,7 +9880,7 @@ impl DiskAPI for LocalDisk {
|
||||
);
|
||||
return Err(e);
|
||||
}
|
||||
// TODO: health check
|
||||
// TODO(backlog): add post-setup disk health verification
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -315,7 +315,7 @@ pub async fn fsync_dir(dir: impl AsRef<Path>) -> io::Result<()> {
|
||||
#[cfg(unix)]
|
||||
{
|
||||
let dir = dir.as_ref().to_path_buf();
|
||||
tokio::task::spawn_blocking(move || fsync_dir_std(dir)).await?
|
||||
fsync_spawn_blocking(move || fsync_dir_std(dir)).await?
|
||||
}
|
||||
|
||||
#[cfg(not(unix))]
|
||||
@@ -683,7 +683,7 @@ async fn fsync_open_dst_dir_group(group: &DstDirFsyncGroup) -> io::Result<()> {
|
||||
#[cfg(test)]
|
||||
let dir = group.dir.clone();
|
||||
let dir_file = group.dir_file.clone();
|
||||
tokio::task::spawn_blocking(move || {
|
||||
fsync_spawn_blocking(move || {
|
||||
#[cfg(test)]
|
||||
{
|
||||
if let Some(kind) = fsync_dir_recorder::take_grouped_failure(&dir) {
|
||||
@@ -1080,6 +1080,44 @@ const TEST_GLOBAL_FILE_SYNCS: usize = 64;
|
||||
|
||||
static FILE_SYNC_PERMITS: LazyLock<Semaphore> = LazyLock::new(|| Semaphore::new(global_file_sync_limit()));
|
||||
static DISK_FILE_SYNC_LIMITERS: LazyLock<Mutex<HashMap<PathBuf, Weak<Semaphore>>>> = LazyLock::new(|| Mutex::new(HashMap::new()));
|
||||
|
||||
/// Dedicated tokio runtime for fsync/fdatasync blocking operations. When
|
||||
/// configured with >1 threads, isolates device-bound fsync from the main
|
||||
/// blocking pool so reads (pread/stat/open) are not starved. `None` means
|
||||
/// fall back to the main runtime (zero behavior change).
|
||||
static FSYNC_RUNTIME: LazyLock<Option<tokio::runtime::Runtime>> = LazyLock::new(|| {
|
||||
let threads =
|
||||
rustfs_utils::get_env_usize(rustfs_config::ENV_FSYNC_BLOCKING_THREADS, rustfs_config::DEFAULT_FSYNC_BLOCKING_THREADS);
|
||||
if threads <= 1 {
|
||||
return None;
|
||||
}
|
||||
let mut builder = tokio::runtime::Builder::new_multi_thread();
|
||||
builder
|
||||
.worker_threads(num_cpus::get().min(8))
|
||||
.max_blocking_threads(threads)
|
||||
.thread_name("rustfs-fsync")
|
||||
.thread_stack_size(512 * 1024)
|
||||
.enable_all();
|
||||
match builder.build() {
|
||||
Ok(rt) => {
|
||||
tracing::info!(threads, "fsync dedicated blocking pool enabled");
|
||||
Some(rt)
|
||||
}
|
||||
Err(err) => {
|
||||
tracing::warn!(%err, "failed to build fsync runtime, falling back to main pool");
|
||||
None
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
/// Spawn a blocking task on the fsync-dedicated runtime if configured,
|
||||
/// otherwise fall back to the main tokio blocking pool.
|
||||
fn fsync_spawn_blocking<T: Send + 'static>(f: impl FnOnce() -> T + Send + 'static) -> tokio::task::JoinHandle<T> {
|
||||
match FSYNC_RUNTIME.as_ref() {
|
||||
Some(rt) => rt.spawn_blocking(f),
|
||||
None => tokio::task::spawn_blocking(f),
|
||||
}
|
||||
}
|
||||
static DISK_VOLUME_MUTATION_LOCKS: LazyLock<Mutex<HashMap<PathBuf, Weak<RwLock<()>>>>> =
|
||||
LazyLock::new(|| Mutex::new(HashMap::new()));
|
||||
type NamespaceMutationLock = AsyncMutex<()>;
|
||||
@@ -1217,7 +1255,7 @@ where
|
||||
F: FnOnce() -> io::Result<T> + Send + 'static,
|
||||
{
|
||||
let (disk_permit, global_permit) = acquire_file_sync_permits(disk_permits).await?;
|
||||
let result = tokio::task::spawn_blocking(move || {
|
||||
let result = fsync_spawn_blocking(move || {
|
||||
let _disk_permit = disk_permit;
|
||||
work()
|
||||
})
|
||||
@@ -2146,7 +2184,7 @@ async fn run_blocking_namespace_file_sync_operation_with_global<T: Send + 'stati
|
||||
wait_started,
|
||||
);
|
||||
let disk_permit = admission.disk_permit.clone();
|
||||
let result = tokio::task::spawn_blocking(move || {
|
||||
let result = fsync_spawn_blocking(move || {
|
||||
let _lease = lease;
|
||||
let _disk_permit = disk_permit;
|
||||
operation()
|
||||
|
||||
@@ -249,7 +249,7 @@ impl PoolEndpointList {
|
||||
endpoint.set_set_index(0);
|
||||
endpoint.set_disk_index(0);
|
||||
|
||||
// TODO Check for cross device mounts if any.
|
||||
// TODO(backlog): check for cross-device mounts in single-drive setup
|
||||
|
||||
return Ok(Self {
|
||||
inner: vec![Endpoints::from(vec![endpoint])],
|
||||
@@ -264,7 +264,7 @@ impl PoolEndpointList {
|
||||
// Convert args to endpoints
|
||||
let mut eps = Endpoints::try_from(set_layout.as_slice())?;
|
||||
|
||||
// TODO Check for cross device mounts if any.
|
||||
// TODO(backlog): check for cross-device mounts in multi-pool setup
|
||||
|
||||
for (disk_idx, ep) in eps.as_mut().iter_mut().enumerate() {
|
||||
ep.set_pool_index(pool_idx);
|
||||
|
||||
@@ -1091,7 +1091,7 @@ impl ObjectInfo {
|
||||
}
|
||||
};
|
||||
|
||||
// TODO:VersionPurgeStatus
|
||||
// TODO(backlog): handle VersionPurgeStatus in object listing
|
||||
let versioned = vcfg.clone().map(|v| v.0.versioned(&entry.name)).unwrap_or_default();
|
||||
objects.push(ObjectInfo::from_file_info(&fi, bucket, &entry.name, versioned));
|
||||
|
||||
|
||||
@@ -1575,7 +1575,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
||||
let parts_metadata = vec![fi.clone(); disks.len()];
|
||||
|
||||
if !user_defined.contains_key("content-type") {
|
||||
// TODO: get content-type
|
||||
// TODO(backlog): detect content-type from part data when header is missing
|
||||
}
|
||||
|
||||
if let Some(sc) = user_defined.get(AMZ_STORAGE_CLASS)
|
||||
@@ -1971,7 +1971,7 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks {
|
||||
return Err(Error::InvalidPart(p.part_num, ext_part.etag.clone(), p.etag.clone().unwrap_or_default()));
|
||||
}
|
||||
|
||||
// TODO: crypto
|
||||
// TODO(backlog): integrate encryption verification during complete multipart
|
||||
|
||||
if (i < uploaded_parts.len() - 1)
|
||||
&& !(opts.data_movement && ext_part.actual_size < 0)
|
||||
|
||||
@@ -6161,7 +6161,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
||||
|
||||
join_all(rollback_futures).await;
|
||||
|
||||
// TODO: add_partial
|
||||
// TODO(backlog): support partial object deletion for multi-part objects
|
||||
|
||||
if let Some(api) = opts.tier_delete_journal_api.as_ref() {
|
||||
for (idx, je) in persisted_journal_entries {
|
||||
@@ -6371,7 +6371,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
||||
}
|
||||
}
|
||||
|
||||
// TODO: Lifecycle
|
||||
// TODO(backlog): integrate lifecycle evaluation before object deletion
|
||||
|
||||
let mut version_found = true;
|
||||
// delete_object_version below derives its own majority quorum from the
|
||||
@@ -6465,7 +6465,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
||||
mark_deleted: mark_delete,
|
||||
mod_time: Some(mod_time),
|
||||
replication_state_internal: opts.delete_replication.as_ref().map(replication_state_to_filemeta),
|
||||
..Default::default() // TODO: Transition
|
||||
..Default::default() // TODO(backlog): populate transition state on delete markers
|
||||
};
|
||||
|
||||
fi.set_tier_free_version_id(&find_vid.to_string());
|
||||
|
||||
@@ -601,7 +601,7 @@ impl ECStore {
|
||||
|
||||
#[instrument(skip(self))]
|
||||
pub(super) async fn handle_list_bucket(&self, opts: &BucketOptions) -> Result<Vec<BucketInfo>> {
|
||||
// TODO: opts.cached
|
||||
// TODO(backlog): support cached bucket listing via opts.cached
|
||||
|
||||
let mut buckets = self.peer_sys.list_bucket(opts).await?;
|
||||
|
||||
|
||||
@@ -4673,7 +4673,7 @@ async fn gather_results(
|
||||
entry.name = entry.name.replace("\\", "/");
|
||||
}
|
||||
|
||||
// TODO: rx.recv()
|
||||
// TODO(backlog): integrate rx.recv() for incremental listing results
|
||||
|
||||
if let Some(marker) = &opts.marker
|
||||
&& ((!opts.include_marker && &entry.name <= marker) || (opts.include_marker && &entry.name < marker))
|
||||
@@ -4703,7 +4703,7 @@ async fn gather_results(
|
||||
continue;
|
||||
}
|
||||
|
||||
// TODO: Lifecycle
|
||||
// TODO(backlog): integrate lifecycle evaluation during object listing
|
||||
|
||||
entries.push(Some(entry));
|
||||
candidate_entries += 1;
|
||||
|
||||
@@ -332,7 +332,7 @@ impl ECStore {
|
||||
let expected_incarnation_id = opts.expected_bucket_incarnation_id;
|
||||
|
||||
if request.prefix.is_empty() {
|
||||
// TODO: return from cache
|
||||
// TODO(backlog): return cached multipart listing when prefix is empty
|
||||
}
|
||||
|
||||
if self.single_pool() {
|
||||
@@ -610,7 +610,7 @@ impl ECStore {
|
||||
let (opts, _bucket_lifecycle_guard) = self.guard_multipart_bucket_incarnation(bucket, opts).await?;
|
||||
let opts = &opts;
|
||||
|
||||
// TODO: defer DeleteUploadID
|
||||
// TODO(backlog): defer DeleteUploadID to background for faster abort response
|
||||
|
||||
if self.single_pool() {
|
||||
return self.pools[0].abort_multipart_upload(bucket, object, upload_id, opts).await;
|
||||
|
||||
@@ -385,7 +385,7 @@ impl ECStore {
|
||||
}
|
||||
|
||||
pub(super) async fn is_suspended(&self, idx: usize) -> bool {
|
||||
// TODO: LOCK
|
||||
// TODO(backlog): acquire pool metadata lock for consistent suspension check
|
||||
|
||||
let pool_meta = self.pool_meta.read().await;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user