fix: resolve pagination fixtures and heal recovery CI failures (#8316)

* test(e2e): bound delimiter pagination fixture uploads

* fix(heal): classify CAS lock failures before publication

* fix(heal): retry unproven absence within the object budget

* docs(heal): clarify object retry classification

* test(s3): seed the delimiter fixture with bounded concurrent PUTs
This commit is contained in:
Chris
2026-10-03 20:40:52 +08:00
committed by GitHub
parent 5efe4c5703
commit 8452910895
9 changed files with 267 additions and 26 deletions
@@ -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<String> = (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<String> = (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;
+4 -1
View File
@@ -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)]
+35
View File
@@ -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<std::time::SystemTime> {
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());
+2 -2
View File
@@ -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
@@ -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);
}
@@ -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() {
+6 -3
View File
@@ -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);
@@ -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);
}
@@ -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'] == '/'