Compare commits

...

4 Commits

Author SHA1 Message Date
cxymds 1a3810991d Merge branch 'main' into cxymds/fix-decommission-listing-failclosed 2026-08-22 11:24:07 +08:00
马登山 bb2330451a perf(ecstore): avoid successful listing name clone 2026-08-22 11:08:04 +08:00
马登山 09fac9a52b fix(ecstore): fail closed on unresolved decommission entries 2026-08-22 11:04:29 +08:00
Zhengchao An 5b951de2b7 test(ci): bound s3-tests failure logs (#6361) 2026-08-22 02:58:01 +00:00
3 changed files with 229 additions and 26 deletions
+201 -24
View File
@@ -36,7 +36,8 @@ use crate::disk::error::DiskError;
use crate::disk::{BUCKET_META_PREFIX, RUSTFS_META_BUCKET}; use crate::disk::{BUCKET_META_PREFIX, RUSTFS_META_BUCKET};
use crate::error::{Error, Result}; use crate::error::{Error, Result};
use crate::error::{ 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::layout::endpoints::EndpointServerPools;
use crate::object_api::{GetObjectReader, ObjectOptions}; use crate::object_api::{GetObjectReader, ObjectOptions};
@@ -773,7 +774,76 @@ async fn load_decommission_entry_exact_versions(
} }
fn resolve_decommission_check_after_list_result(list_result: Result<()>, entry_error: Option<Error>) -> Result<()> { 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,
))
}
async fn record_decommission_entry_error(
entry_error: &Arc<tokio::sync::Mutex<Option<Error>>>,
rx: &CancellationToken,
err: Error,
) {
if rx.is_cancelled() {
return;
}
let mut first_err = entry_error.lock().await;
if first_err.is_none() && !rx.is_cancelled() {
*first_err = Some(err);
rx.cancel();
}
} }
fn resolve_decommission_pool_meta_reload_result(result: Result<()>, stage: &str) -> Result<()> { fn resolve_decommission_pool_meta_reload_result(result: Result<()>, stage: &str) -> Result<()> {
@@ -3538,6 +3608,7 @@ impl ECStore {
let rx_clone = rx.clone(); let rx_clone = rx.clone();
let bi = bi.clone(); let bi = bi.clone();
let set_id = set_idx; let set_id = set_idx;
let listing_entry_error = entry_error.clone();
let worker = tokio::spawn(async move { let worker = tokio::spawn(async move {
let _listing_permit = listing_permit; let _listing_permit = listing_permit;
run_decommission_listing_with_retry( run_decommission_listing_with_retry(
@@ -3551,7 +3622,11 @@ impl ECStore {
let set = set.clone(); let set = set.clone();
let rx = rx_clone.clone(); let rx = rx_clone.clone();
let bucket = bi.clone(); let bucket = bi.clone();
async move { set.list_objects_to_decommission(rx, bucket, callback).await } let entry_error = listing_entry_error.clone();
async move {
set.list_objects_to_decommission(rx, bucket, callback, entry_error.clone(), idx, set_id)
.await
}
}, },
) )
.await .await
@@ -3581,11 +3656,7 @@ impl ECStore {
wait_decommission_worker_drain(&workers, worker_limit).await?; wait_decommission_worker_drain(&workers, worker_limit).await?;
if let Some(err) = listing_worker_error { if let Some(err) = resolve_decommission_listing_error(listing_worker_error, entry_error.lock().await.clone()) {
return Err(err);
}
if let Some(err) = entry_error.lock().await.clone() {
return Err(err); return Err(err);
} }
@@ -4191,7 +4262,7 @@ impl ECStore {
let buckets = self.get_buckets_to_decommission().await?; let buckets = self.get_buckets_to_decommission().await?;
let pool = self.pools[idx].clone(); 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 { for bucket_info in &buckets {
let mut lifecycle_config = None; let mut lifecycle_config = None;
let mut object_lock_config = None; let mut object_lock_config = None;
@@ -4286,7 +4357,7 @@ impl ECStore {
}); });
let list_result = set 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; .await;
let entry_error = entry_error.lock().await.clone(); let entry_error = entry_error.lock().await.clone();
resolve_decommission_check_after_list_result(list_result, entry_error)?; resolve_decommission_check_after_list_result(list_result, entry_error)?;
@@ -5021,12 +5092,15 @@ mod tests {
pub type ListCallback = Arc<dyn Fn(MetaCacheEntry) -> BoxFuture<'static, ()> + Send + Sync + 'static>; pub type ListCallback = Arc<dyn Fn(MetaCacheEntry) -> BoxFuture<'static, ()> + Send + Sync + 'static>;
impl SetDisks { impl SetDisks {
#[tracing::instrument(skip(self, rx, cb_func))] #[tracing::instrument(skip(self, rx, cb_func, entry_error))]
async fn list_objects_to_decommission( async fn list_objects_to_decommission(
self: &Arc<Self>, self: &Arc<Self>,
rx: CancellationToken, rx: CancellationToken,
bucket_info: DecomBucketInfo, bucket_info: DecomBucketInfo,
cb_func: ListCallback, cb_func: ListCallback,
entry_error: Arc<tokio::sync::Mutex<Option<Error>>>,
pool_index: usize,
set_index: usize,
) -> Result<()> { ) -> Result<()> {
let (disks, _) = self.get_online_disks_with_healing(false).await; let (disks, _) = self.get_online_disks_with_healing(false).await;
ensure_decommission_listing_disks_available(!disks.is_empty(), &bucket_info.name)?; ensure_decommission_listing_disks_available(!disks.is_empty(), &bucket_info.name)?;
@@ -5041,6 +5115,12 @@ impl SetDisks {
}; };
let cb1 = cb_func.clone(); 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( list_path_raw(
rx, rx,
@@ -5053,20 +5133,51 @@ impl SetDisks {
skip_walkdir_total_timeout: true, skip_walkdir_total_timeout: true,
walkdir_stall_timeout: Some(DECOMMISSION_BACKGROUND_WALKDIR_STALL_TIMEOUT), walkdir_stall_timeout: Some(DECOMMISSION_BACKGROUND_WALKDIR_STALL_TIMEOUT),
agreed: Some(Box::new(move |entry: MetaCacheEntry| Box::pin(cb1(entry)))), 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 resolver = resolver.clone();
let cb_func = cb_func.clone(); let cb_func = cb_func.clone();
match entries.resolve(resolver) { let bucket = unresolved_bucket.clone();
Some(entry) => { 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); warn!("decommission_pool: list_objects_to_decommission get {}", &entry.name);
Box::pin(async move { Box::pin(async move {
cb_func(entry).await; cb_func(entry).await;
}) })
} }
None => { Err(err) => Box::pin(async move {
warn!("decommission_pool: list_objects_to_decommission get none"); if unresolved_rx.is_cancelled() {
Box::pin(async {}) 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() ..Default::default()
@@ -5074,6 +5185,10 @@ impl SetDisks {
) )
.await?; .await?;
if let Some(err) = entry_error.lock().await.clone() {
return Err(err);
}
Ok(()) Ok(())
} }
} }
@@ -5279,11 +5394,12 @@ mod pools_tests {
has_active_decommission_canceler, is_decommission_active, is_decommission_cancel_requested, has_active_decommission_canceler, is_decommission_active, is_decommission_cancel_requested,
load_decommission_entry_versions, local_decommission_queue_prefix, mark_decommission_bucket_done, 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, 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, pool_meta_has_active_decommission, record_decommission_entry_error, require_decommission_store,
resolve_decommission_bucket_state, resolve_decommission_check_after_list_result, resolve_decommission_bucket_done_save_result, resolve_decommission_bucket_state,
resolve_decommission_entry_cleanup_delete_result, resolve_decommission_entry_exact_versions, resolve_decommission_check_after_list_result, resolve_decommission_entry_cleanup_delete_result,
resolve_decommission_entry_reload_result, resolve_decommission_listing_worker_result, resolve_decommission_entry_exact_versions, resolve_decommission_entry_reload_result, resolve_decommission_listing_error,
resolve_decommission_optional_bucket_config_result, resolve_decommission_pool_meta_reload_result, 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_preflight_heal_result, resolve_decommission_progress_save_result,
resolve_decommission_spawn_failure_result, resolve_decommission_terminal_mark_after_error_result, resolve_decommission_spawn_failure_result, resolve_decommission_terminal_mark_after_error_result,
resolve_decommission_terminal_mark_result, resolve_decommission_update_after_result, resolve_decommission_terminal_mark_result, resolve_decommission_update_after_result,
@@ -5302,7 +5418,9 @@ mod pools_tests {
use crate::error::{Error, StorageError}; use crate::error::{Error, StorageError};
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints}; use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
use crate::services::rebalance::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats}; 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 rustfs_rio::Index;
use std::sync::{ use std::sync::{
Arc, Arc,
@@ -6321,6 +6439,65 @@ mod pools_tests {
assert!(matches!(err, Error::SlowDown)); 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] #[test]
fn test_resolve_decommission_check_after_list_result_returns_list_result_without_entry_error() { 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) let err = resolve_decommission_check_after_list_result(Err(Error::OperationCanceled), None)
+26 -1
View File
@@ -206,6 +206,13 @@ def check_runner_selection(root: Path) -> list[str]:
return errors return errors
def check_s3_tests_runner(root: Path) -> list[str]:
runner = (root / "scripts/s3-tests/run.sh").read_text()
if "--showlocals" in runner:
return ["scripts/s3-tests/run.sh: pytest failure diagnostics must not dump local values"]
return []
def profile_selection(root: Path, profile: str) -> str: def profile_selection(root: Path, profile: str) -> str:
if not re.fullmatch(r"e2e-[a-z0-9-]+", profile): if not re.fullmatch(r"e2e-[a-z0-9-]+", profile):
raise ValueError(f"invalid e2e profile name: {profile}") raise ValueError(f"invalid e2e profile name: {profile}")
@@ -272,6 +279,7 @@ def validate(root: Path) -> list[str]:
errors.extend(check_e2e_modules(root)) errors.extend(check_e2e_modules(root))
errors.extend(check_fuzz_targets(root)) errors.extend(check_fuzz_targets(root))
errors.extend(check_runner_selection(root)) errors.extend(check_runner_selection(root))
errors.extend(check_s3_tests_runner(root))
errors.extend(check_profile_definitions(root)) errors.extend(check_profile_definitions(root))
return errors return errors
@@ -341,6 +349,23 @@ class SelfTests(unittest.TestCase):
) )
self.assertEqual(len(check_fuzz_targets(root)), 1) self.assertEqual(len(check_fuzz_targets(root)), 1)
def test_s3_runner_rejects_unbounded_failure_locals(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
runner = root / "scripts/s3-tests/run.sh"
runner.parent.mkdir(parents=True)
runner.write_text("tox -- -vv -ra --tb=long\n")
self.assertEqual(check_s3_tests_runner(root), [])
runner.write_text("tox -- -vv -ra --showlocals --tb=long\n")
self.assertEqual(len(check_s3_tests_runner(root)), 1)
with (
mock.patch(__name__ + ".check_e2e_modules", return_value=[]),
mock.patch(__name__ + ".check_fuzz_targets", return_value=[]),
mock.patch(__name__ + ".check_runner_selection", return_value=[]),
mock.patch(__name__ + ".check_profile_definitions", return_value=[]),
):
self.assertEqual(len(validate(root)), 1)
def test_profile_listing_enforces_selection(self) -> None: def test_profile_listing_enforces_selection(self) -> None:
with tempfile.TemporaryDirectory() as tmp: with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp) root = Path(tmp)
@@ -411,7 +436,7 @@ def main() -> int:
for error in errors: for error in errors:
print(f"ERROR: {error}", file=sys.stderr) print(f"ERROR: {error}", file=sys.stderr)
return 1 return 1
print("OK: e2e modules, runner selection, fuzz matrices, and profile guards are wired") print("OK: e2e modules, runner selection, fuzz matrices, profiles, and bounded diagnostics are wired")
return 0 return 0
+2 -1
View File
@@ -1028,10 +1028,11 @@ else
fi fi
# Run tests from s3tests/functional # Run tests from s3tests/functional
# Failure locals can contain multi-MiB request bodies; keep tracebacks without expanding local values.
set +e set +e
S3TEST_CONF="${CONF_OUTPUT_PATH}" \ S3TEST_CONF="${CONF_OUTPUT_PATH}" \
tox -- \ tox -- \
-vv -ra --showlocals --tb=long \ -vv -ra --tb=long \
--maxfail="${MAXFAIL}" \ --maxfail="${MAXFAIL}" \
--timeout="${TEST_TIMEOUT}" \ --timeout="${TEST_TIMEOUT}" \
--junitxml="${ARTIFACTS_DIR}/junit.xml" \ --junitxml="${ARTIFACTS_DIR}/junit.xml" \