diff --git a/crates/e2e_test/src/list_objects_v2_pagination_test.rs b/crates/e2e_test/src/list_objects_v2_pagination_test.rs index cf2fcea2a..bc2ba8090 100644 --- a/crates/e2e_test/src/list_objects_v2_pagination_test.rs +++ b/crates/e2e_test/src/list_objects_v2_pagination_test.rs @@ -811,23 +811,25 @@ mod tests { create_bucket(&client, bucket).await.expect("Failed to create bucket"); - // Create 1000 objects: 100 directories with 10 files each - let mut all_keys = Vec::new(); - let dirs: Vec = (0..100).map(|i| format!("dir-{:03}/", i)).collect(); - for dir in &dirs { - for i in 0..10 { - let key = format!("{dir}file{:02}.txt", i); - client - .put_object() - .bucket(bucket) - .key(&key) - .body(ByteStream::from_static(b"x")) - .send() - .await - .expect("Failed to put object"); - all_keys.push(key); - } - } + // Keep every fixture key while overlapping durable PUTs within a bounded fanout. + let all_keys: Vec = (0..100) + .flat_map(|dir| (0..10).map(move |file| format!("dir-{dir:03}/file{file:02}.txt"))) + .collect(); + stream::iter(&all_keys) + .for_each_concurrent(16, |key| { + let client = &client; + async move { + client + .put_object() + .bucket(bucket) + .key(key) + .body(ByteStream::from_static(b"x")) + .send() + .await + .unwrap_or_else(|err| panic!("Failed to put fixture object {key}: {err}")); + } + }) + .await; eprintln!("Seeded {} objects in {bucket}; starting ListObjectsV2 pagination", all_keys.len()); let deadline = Instant::now() + PAGINATION_TIMEOUT; diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 5a3994250..124799644 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -9499,7 +9499,9 @@ impl DiskAPI for LocalDisk { .write(true) .open(&lock_path) .map_err(DiskError::conditional_file_not_committed)?; - flock(&lock, FlockOperation::NonBlockingLockExclusive).map_err(std::io::Error::from)?; + flock(&lock, FlockOperation::NonBlockingLockExclusive) + .map_err(std::io::Error::from) + .map_err(DiskError::conditional_file_not_committed)?; let result = (|| { let current = match std::fs::read(&file_path) { Ok(current) => Some(current), @@ -25138,6 +25140,7 @@ mod test { .expect_err("contended conditional update must retry"); assert!(matches!(err, DiskError::Io(ref err) if err.kind() == ErrorKind::WouldBlock)); + assert!(err.is_conditional_file_not_committed(), "lock contention cannot publish target bytes"); } #[cfg(windows)] diff --git a/crates/heal/src/error.rs b/crates/heal/src/error.rs index 8d2f32aae..a771a6663 100644 --- a/crates/heal/src/error.rs +++ b/crates/heal/src/error.rs @@ -177,6 +177,21 @@ impl Error { } } + /// Retry an unproven missing object only within the per-object budget. + /// Generic task retries retain the narrower classification above. + pub(crate) fn is_recoverable_object_heal(&self) -> bool { + self.is_recoverable_heal() + || matches!( + self, + Error::Storage( + EcstoreError::FileNotFound + | EcstoreError::FileVersionNotFound + | EcstoreError::ObjectNotFound(_, _) + | EcstoreError::VersionNotFound(_, _, _) + ) | Error::Disk(DiskError::FileNotFound | DiskError::FileVersionNotFound) + ) + } + pub(crate) fn dangling_delete_retry_not_before(&self) -> Option { let after = match self { Self::Storage(error) => error.dangling_delete_retry_after(), @@ -395,6 +410,26 @@ mod tests { assert!(Error::Storage(EcstoreError::FaultyRemoteDisk).is_recoverable_heal()); } + #[test] + fn typed_missing_objects_retry_only_within_the_object_budget() { + for error in [ + Error::Storage(EcstoreError::FileNotFound), + Error::Storage(EcstoreError::FileVersionNotFound), + Error::Storage(EcstoreError::ObjectNotFound("bucket".to_owned(), "object".to_owned())), + Error::Storage(EcstoreError::VersionNotFound( + "bucket".to_owned(), + "object".to_owned(), + "version".to_owned(), + )), + Error::Disk(DiskError::FileNotFound), + Error::Disk(DiskError::FileVersionNotFound), + ] { + assert!(error.is_recoverable_object_heal(), "{error:?}"); + assert!(!error.is_recoverable_heal(), "missing objects must not requeue the whole task: {error:?}"); + } + assert!(!Error::other("file not found").is_recoverable_object_heal()); + } + #[test] fn task_timeout_is_terminal() { assert!(!Error::TaskTimeout.is_recoverable_heal()); diff --git a/crates/heal/src/heal/erasure_healer/admin.rs b/crates/heal/src/heal/erasure_healer/admin.rs index f7fb6d88a..5c9a686cd 100644 --- a/crates/heal/src/heal/erasure_healer/admin.rs +++ b/crates/heal/src/heal/erasure_healer/admin.rs @@ -196,7 +196,7 @@ pub(super) async fn heal_object( return ((size, Err(error.unwrap_or(Error::TaskCancelled))), None); } if let Some(error) = error.as_ref() - && error.is_recoverable_heal() + && error.is_recoverable_object_heal() && failures < MAX_BUCKET_OBJECT_HEAL_RETRIES { failures += 1; @@ -224,7 +224,7 @@ pub(super) async fn heal_object( None => Ok(!matches!(disposition, HealObjectDisposition::AuthoritativelyAbsent)), Some(error) => { failures += 1; - if error.is_recoverable_heal() { + if error.is_recoverable_object_heal() { disposition = HealObjectDisposition::Deferred { reason: if error.is_dangling_delete_grace() { HealDeferredReason::DanglingDeleteGrace diff --git a/crates/heal/src/heal/erasure_healer/resume_loop_tests/admin_outcome.rs b/crates/heal/src/heal/erasure_healer/resume_loop_tests/admin_outcome.rs index 0590acb2b..e893f8eca 100644 --- a/crates/heal/src/heal/erasure_healer/resume_loop_tests/admin_outcome.rs +++ b/crates/heal/src/heal/erasure_healer/resume_loop_tests/admin_outcome.rs @@ -335,9 +335,23 @@ async fn admin_erasure_dispositions_do_not_invent_success_and_retries_count_once .expect_err("unhealed objects must retain the failed task status"); let outcome = task.get_outcome().await; let c = &outcome.counters; - assert_eq!((c.processed, c.healed, c.unchanged, c.failed, c.skipped, c.unknown), (5, 0, 2, 1, 2, 1)); - assert_eq!(c.attempt_failures, 5); + assert_eq!((c.processed, c.healed, c.unchanged, c.failed, c.skipped, c.unknown), (5, 0, 2, 0, 3, 1)); + assert_eq!(c.attempt_failures, 8); + assert_eq!(storage.calls().iter().filter(|(name, _)| name == "gone").count(), 4); assert_eq!(storage.calls().iter().filter(|(name, _)| name == "offline").count(), 4); + let gone = outcome + .objects + .iter() + .find(|item| item.identity.object == "gone") + .expect("missing version outcome"); + assert_eq!( + gone.disposition, + HealObjectDisposition::Deferred { + reason: HealDeferredReason::TransientExistenceCheck, + retry_not_before: None, + }, + "exhausting retries cannot certify the missing version as absent" + ); assert_partition(&outcome); } diff --git a/crates/heal/src/heal/manager/tests/root_recovery.rs b/crates/heal/src/heal/manager/tests/root_recovery.rs index d792064cb..732c83383 100644 --- a/crates/heal/src/heal/manager/tests/root_recovery.rs +++ b/crates/heal/src/heal/manager/tests/root_recovery.rs @@ -728,6 +728,90 @@ async fn root_recovery_orphan_report_never_commits_or_retires_pending_work() { assert_eq!(manager.root_recovery.pending().await.expect("pending survives GC").len(), 1); } +#[cfg(unix)] +#[tokio::test] +async fn root_recovery_lock_contention_preserves_a_single_durable_owner() { + let (_first_temp, first_disk) = recovery_disk().await; + let (_second_temp, second_disk) = recovery_disk().await; + let (first_disk, second_disk) = ordered_recovery_disks(first_disk, second_disk); + let lock_owner = |disk: &DiskStore| { + let path = std::path::Path::new(&disk.endpoint().to_string()) + .join(RUSTFS_META_BUCKET) + .join(".rustfs-cas.lock"); + let file = std::fs::OpenOptions::new() + .create(true) + .truncate(false) + .read(true) + .write(true) + .open(path) + .expect("open control-file lock"); + file.lock().expect("hold control-file lock"); + file + }; + let first_lock = lock_owner(&first_disk); + let manager = recovery_manager(vec![first_disk.clone(), second_disk.clone()]); + let mut request = admin_request(HealType::Object { + bucket: "bucket".to_string(), + object: "object".to_string(), + version_id: None, + }); + let receipt = manager + .submit_heal_request_with_receipt(request.clone()) + .await + .expect("uncontended disk should own the durable heal intent"); + assert_eq!(receipt.result, HealAdmissionResult::Accepted); + let path = format!("root-heal-{}.json", request.id); + assert!(matches!( + first_disk.read_all(RUSTFS_META_BUCKET, &path).await, + Err(DiskError::FileNotFound) + )); + let original = second_disk.read_all(RUSTFS_META_BUCKET, &path).await.expect("durable intent"); + + drop(first_lock); + let _second_lock = lock_owner(&second_disk); + request.retry_attempts = 1; + manager + .root_recovery + .persist(&request) + .await + .expect_err("existing ownership must not migrate when its lock is contended"); + assert!(matches!( + first_disk.read_all(RUSTFS_META_BUCKET, &path).await, + Err(DiskError::FileNotFound) + )); + assert_eq!( + second_disk + .read_all(RUSTFS_META_BUCKET, &path) + .await + .expect("unchanged owner"), + original + ); + let pending = manager.root_recovery.pending().await.expect("single durable owner"); + assert_eq!(pending.len(), 1); + assert_eq!(pending[0].id, request.id); + assert_eq!(pending[0].retry_attempts, 0); + + let _first_lock = lock_owner(&first_disk); + let blocked = admin_request(HealType::Object { + bucket: "bucket".to_string(), + object: "blocked".to_string(), + version_id: None, + }); + manager + .root_recovery + .persist(&blocked) + .await + .expect_err("all owners remain contended"); + let blocked_path = format!("root-heal-{}.json", blocked.id); + for disk in [&first_disk, &second_disk] { + assert!(matches!( + disk.read_all(RUSTFS_META_BUCKET, &blocked_path).await, + Err(DiskError::FileNotFound) + )); + } + assert_eq!(manager.root_recovery.pending().await.expect("original intent remains").len(), 1); +} + #[cfg(unix)] #[tokio::test] async fn root_recovery_new_intent_skips_prepublication_read_only_owner() { diff --git a/crates/heal/src/heal/task/heal_bucket.rs b/crates/heal/src/heal/task/heal_bucket.rs index 8abd2b4f8..73bc20291 100644 --- a/crates/heal/src/heal/task/heal_bucket.rs +++ b/crates/heal/src/heal/task/heal_bucket.rs @@ -762,7 +762,10 @@ impl HealTask { error = %err, "Heal bucket object repair skipped due to transient metadata error" ); - } else if !age_exhausted && err.is_recoverable_heal() && retry_attempt < MAX_BUCKET_OBJECT_HEAL_RETRIES { + } else if !age_exhausted + && err.is_recoverable_object_heal() + && retry_attempt < MAX_BUCKET_OBJECT_HEAL_RETRIES + { terminal_outcome = false; debug!( target: "rustfs::heal::task", @@ -782,13 +785,13 @@ impl HealTask { inline_retry = Some(item); } } else { - disposition = HealObjectDisposition::Failed(if age_exhausted || err.is_recoverable_heal() { + disposition = HealObjectDisposition::Failed(if age_exhausted || err.is_recoverable_object_heal() { HealFailureClass::RetryExhausted } else { HealFailureClass::Permanent }); telemetry_unknown |= !increment_counter(&mut failed); - if age_exhausted || err.is_recoverable_heal() { + if age_exhausted || err.is_recoverable_object_heal() { retryable_failed = retryable_failed.saturating_add(1); } else { permanent_failed = permanent_failed.saturating_add(1); diff --git a/crates/heal/src/heal/task/tests/concurrent_delete.rs b/crates/heal/src/heal/task/tests/concurrent_delete.rs index 5fbac9047..485e6fd81 100644 --- a/crates/heal/src/heal/task/tests/concurrent_delete.rs +++ b/crates/heal/src/heal/task/tests/concurrent_delete.rs @@ -123,7 +123,7 @@ async fn admin_dry_run_can_observe_without_bucket_identity() { #[tokio::test(start_paused = true)] async fn admin_traversal_never_converts_unproven_errors_to_absence() { for (error, class) in [ - (MockHealObjectOutcome::MissingVersion, HealFailureClass::Permanent), + (MockHealObjectOutcome::MissingVersion, HealFailureClass::RetryExhausted), (MockHealObjectOutcome::PermissionDenied, HealFailureClass::Permanent), (MockHealObjectOutcome::ErrOther("file not found"), HealFailureClass::Permanent), (MockHealObjectOutcome::RetryableReadQuorum, HealFailureClass::RetryExhausted), @@ -147,3 +147,73 @@ async fn admin_traversal_never_converts_unproven_errors_to_absence() { assert_eq!(outcome.objects[0].disposition, HealObjectDisposition::Failed(class)); } } + +#[tokio::test(start_paused = true)] +async fn admin_traversal_retries_unproven_absence_until_storage_certifies_it() { + let incarnation = Uuid::new_v4(); + let receipt = object_receipt("object-a", None, HealObjectDisposition::AuthoritativelyAbsent, incarnation); + let storage = Arc::new(MockStorage { + bucket_incarnation_id: Mutex::new(Some(incarnation)), + heal_object_outcomes: Mutex::new(HashMap::from([( + "object-a".to_string(), + VecDeque::from([MockHealObjectOutcome::MissingVersion]), + )])), + heal_object_receipts: Mutex::new(HashMap::from([("object-a".to_string(), VecDeque::from([receipt.clone(), receipt]))])), + ..Default::default() + }); + let task = admin_traversal(HealType::Cluster, storage.clone(), false); + task.execute() + .await + .expect("unproven absence must use the existing retry budget"); + let outcome = task.get_outcome().await; + let object = outcome + .objects + .iter() + .find(|item| item.identity.object == "object-a") + .expect("object outcome"); + assert_eq!(object.disposition, HealObjectDisposition::AuthoritativelyAbsent); + assert_eq!(object.identity.bucket_incarnation_id, Some(incarnation)); + assert_eq!(outcome.counters.attempt_failures, 1, "the errored receipt cannot certify absence"); + assert_eq!(outcome.counters.failed, 0); + assert_eq!( + storage.heal_object_calls.lock().expect("object calls").as_slice(), + ["object-a", "object-b", "object-a"] + ); +} + +#[tokio::test(start_paused = true)] +async fn admin_traversal_unproven_absence_exhausts_retries_without_positive_outcome() { + let incarnation = Uuid::new_v4(); + let attempts = MAX_BUCKET_OBJECT_HEAL_RETRIES + 1; + let storage = Arc::new(MockStorage { + bucket_incarnation_id: Mutex::new(Some(incarnation)), + heal_object_outcomes: Mutex::new(HashMap::from([( + "object-a".to_string(), + (0..attempts).map(|_| MockHealObjectOutcome::MissingVersion).collect(), + )])), + heal_object_receipts: Mutex::new(HashMap::from([( + "object-a".to_string(), + (0..attempts) + .map(|_| object_receipt("object-a", None, HealObjectDisposition::AuthoritativelyAbsent, incarnation)) + .collect(), + )])), + ..Default::default() + }); + let task = admin_traversal(HealType::Cluster, storage.clone(), false); + task.execute().await.expect_err("unproven absence cannot finish successfully"); + let outcome = task.get_outcome().await; + let object = outcome + .objects + .iter() + .find(|item| item.identity.object == "object-a") + .expect("object outcome"); + assert_eq!(object.disposition, HealObjectDisposition::Failed(HealFailureClass::RetryExhausted)); + assert_eq!(outcome.counters.attempt_failures, u64::from(attempts)); + assert_eq!((outcome.counters.failed, outcome.counters.healed), (1, 0)); + let calls = storage.heal_object_calls.lock().expect("object calls"); + assert_eq!( + calls.iter().filter(|object| object.as_str() == "object-a").count(), + usize::try_from(attempts).expect("attempts fit") + ); + assert_eq!(calls.iter().filter(|object| object.as_str() == "object-b").count(), 1); +} diff --git a/scripts/s3-tests/patches/0005-delimiter-fixture-concurrent-put.patch b/scripts/s3-tests/patches/0005-delimiter-fixture-concurrent-put.patch new file mode 100644 index 000000000..49d3cbfe5 --- /dev/null +++ b/scripts/s3-tests/patches/0005-delimiter-fixture-concurrent-put.patch @@ -0,0 +1,30 @@ +diff --git a/s3tests/functional/test_s3.py b/s3tests/functional/test_s3.py +--- a/s3tests/functional/test_s3.py ++++ b/s3tests/functional/test_s3.py +@@ -11,6 +11,7 @@ + import re + import pytz + from collections import OrderedDict ++from concurrent.futures import ThreadPoolExecutor + import requests + import json + import base64 +@@ -687,8 +688,16 @@ + key_names = ['0/'] + ['0/%s' % i for i in range(1000, 1999)] + key_names2 = ['1999', '1999#', '1999+', '2000'] + key_names += key_names2 +- bucket_name = _create_objects(keys=key_names) +- client = get_client() ++ bucket_name = get_new_bucket_name() ++ get_new_bucket_resource(name=bucket_name) ++ client = get_client() ++ ++ # Share the low-level client; Boto3 resources are not thread-safe. ++ with ThreadPoolExecutor(max_workers=16) as executor: ++ uploads = [executor.submit(client.put_object, Bucket=bucket_name, ++ Body=key, Key=key) for key in key_names] ++ for upload in uploads: ++ upload.result() + + response = client.list_objects(Bucket=bucket_name, Delimiter='/') + assert response['Delimiter'] == '/'