Compare commits

..

1 Commits

Author SHA1 Message Date
cxymds f9d45e41e1 fix(quota): account compressed deletes by committed size (#6365) 2026-08-22 11:39:01 +00:00
33 changed files with 833 additions and 938 deletions
-56
View File
@@ -901,10 +901,6 @@ pub struct Metrics {
scanner_cycle_max_duration_millis: AtomicU64, scanner_cycle_max_duration_millis: AtomicU64,
scanner_cycle_max_objects: AtomicU64, scanner_cycle_max_objects: AtomicU64,
scanner_cycle_max_directories: AtomicU64, scanner_cycle_max_directories: AtomicU64,
scanner_cycle_timeout_total: AtomicU64,
scanner_cycle_recovery_required_total: AtomicU64,
scanner_cycle_last_progress_age_seconds: AtomicU64,
scanner_leader_lease_without_progress: AtomicBool,
scanner_bitrot_cycle_enabled: AtomicBool, scanner_bitrot_cycle_enabled: AtomicBool,
scanner_bitrot_cycle_millis: AtomicU64, scanner_bitrot_cycle_millis: AtomicU64,
scanner_checkpoint: Mutex<Option<ScannerCheckpointReport>>, scanner_checkpoint: Mutex<Option<ScannerCheckpointReport>>,
@@ -1374,14 +1370,6 @@ pub struct ScannerMetricsReport {
#[serde(default)] #[serde(default)]
pub cycle_max_directories: u64, pub cycle_max_directories: u64,
#[serde(default)] #[serde(default)]
pub cycle_timeout_total: u64,
#[serde(default)]
pub cycle_recovery_required_total: u64,
#[serde(default)]
pub cycle_last_progress_age: u64,
#[serde(default)]
pub leader_lease_without_progress: bool,
#[serde(default)]
pub bitrot_cycle_enabled: bool, pub bitrot_cycle_enabled: bool,
#[serde(default)] #[serde(default)]
pub bitrot_cycle_seconds: f64, pub bitrot_cycle_seconds: f64,
@@ -1442,9 +1430,6 @@ const OTEL_SCANNER_BUCKETS_SCANNED: &str = "rustfs_scanner_buckets_scanned_total
const OTEL_SCANNER_CYCLES: &str = "rustfs_scanner_cycles_total"; const OTEL_SCANNER_CYCLES: &str = "rustfs_scanner_cycles_total";
const OTEL_SCANNER_CYCLE_DURATION_SECONDS: &str = "rustfs_scanner_cycle_duration_seconds"; const OTEL_SCANNER_CYCLE_DURATION_SECONDS: &str = "rustfs_scanner_cycle_duration_seconds";
const OTEL_SCANNER_BUCKET_DRIVE_DURATION_SECONDS: &str = "rustfs_scanner_bucket_drive_duration_seconds"; const OTEL_SCANNER_BUCKET_DRIVE_DURATION_SECONDS: &str = "rustfs_scanner_bucket_drive_duration_seconds";
const OTEL_SCANNER_CYCLE_TIMEOUT_TOTAL: &str = "rustfs_scanner_cycle_timeout_total";
const OTEL_SCANNER_CYCLE_LAST_PROGRESS_AGE: &str = "rustfs_scanner_cycle_last_progress_age";
const OTEL_SCANNER_LEADER_LEASE_WITHOUT_PROGRESS: &str = "rustfs_scanner_leader_lease_without_progress";
fn scan_cycle_result_label(result: u8) -> &'static str { fn scan_cycle_result_label(result: u8) -> &'static str {
match result { match result {
@@ -1928,10 +1913,6 @@ impl Metrics {
scanner_cycle_max_duration_millis: AtomicU64::new(0), scanner_cycle_max_duration_millis: AtomicU64::new(0),
scanner_cycle_max_objects: AtomicU64::new(0), scanner_cycle_max_objects: AtomicU64::new(0),
scanner_cycle_max_directories: AtomicU64::new(0), scanner_cycle_max_directories: AtomicU64::new(0),
scanner_cycle_timeout_total: AtomicU64::new(0),
scanner_cycle_recovery_required_total: AtomicU64::new(0),
scanner_cycle_last_progress_age_seconds: AtomicU64::new(0),
scanner_leader_lease_without_progress: AtomicBool::new(false),
scanner_bitrot_cycle_enabled: AtomicBool::new(false), scanner_bitrot_cycle_enabled: AtomicBool::new(false),
scanner_bitrot_cycle_millis: AtomicU64::new(0), scanner_bitrot_cycle_millis: AtomicU64::new(0),
scanner_checkpoint: Mutex::new(None), scanner_checkpoint: Mutex::new(None),
@@ -2431,29 +2412,12 @@ impl Metrics {
.store(cycle_max_objects.unwrap_or_default(), Ordering::Relaxed); .store(cycle_max_objects.unwrap_or_default(), Ordering::Relaxed);
self.scanner_cycle_max_directories self.scanner_cycle_max_directories
.store(cycle_max_directories.unwrap_or_default(), Ordering::Relaxed); .store(cycle_max_directories.unwrap_or_default(), Ordering::Relaxed);
self.scanner_leader_lease_without_progress.store(false, Ordering::Relaxed);
self.scanner_cycle_last_progress_age_seconds.store(0, Ordering::Relaxed);
metrics::gauge!(OTEL_SCANNER_LEADER_LEASE_WITHOUT_PROGRESS).set(0.0);
metrics::gauge!(OTEL_SCANNER_CYCLE_LAST_PROGRESS_AGE).set(0.0);
self.scanner_bitrot_cycle_enabled self.scanner_bitrot_cycle_enabled
.store(bitrot_cycle.is_some(), Ordering::Relaxed); .store(bitrot_cycle.is_some(), Ordering::Relaxed);
self.scanner_bitrot_cycle_millis self.scanner_bitrot_cycle_millis
.store(bitrot_cycle.map(duration_millis_saturated).unwrap_or_default(), Ordering::Relaxed); .store(bitrot_cycle.map(duration_millis_saturated).unwrap_or_default(), Ordering::Relaxed);
} }
pub fn record_scanner_cycle_timeout(&self, recovery_required: bool, progress_age: Duration) {
self.scanner_cycle_timeout_total.fetch_add(1, Ordering::Relaxed);
if recovery_required {
self.scanner_cycle_recovery_required_total.fetch_add(1, Ordering::Relaxed);
}
self.scanner_cycle_last_progress_age_seconds
.store(progress_age.as_secs(), Ordering::Relaxed);
self.scanner_leader_lease_without_progress.store(true, Ordering::Relaxed);
metrics::counter!(OTEL_SCANNER_CYCLE_TIMEOUT_TOTAL).increment(1);
metrics::gauge!(OTEL_SCANNER_CYCLE_LAST_PROGRESS_AGE).set(progress_age.as_secs_f64());
metrics::gauge!(OTEL_SCANNER_LEADER_LEASE_WITHOUT_PROGRESS).set(1.0);
}
pub fn record_scanner_set_scan_state(&self, concurrency_limit: Option<usize>, queued: Option<usize>, active: Option<usize>) { pub fn record_scanner_set_scan_state(&self, concurrency_limit: Option<usize>, queued: Option<usize>, active: Option<usize>) {
if let Some(concurrency_limit) = concurrency_limit { if let Some(concurrency_limit) = concurrency_limit {
self.scanner_set_scan_concurrency_limit self.scanner_set_scan_concurrency_limit
@@ -3301,10 +3265,6 @@ impl Metrics {
m.cycle_max_duration_seconds = self.scanner_cycle_max_duration_millis.load(Ordering::Relaxed) as f64 / 1000.0; m.cycle_max_duration_seconds = self.scanner_cycle_max_duration_millis.load(Ordering::Relaxed) as f64 / 1000.0;
m.cycle_max_objects = self.scanner_cycle_max_objects.load(Ordering::Relaxed); m.cycle_max_objects = self.scanner_cycle_max_objects.load(Ordering::Relaxed);
m.cycle_max_directories = self.scanner_cycle_max_directories.load(Ordering::Relaxed); m.cycle_max_directories = self.scanner_cycle_max_directories.load(Ordering::Relaxed);
m.cycle_timeout_total = self.scanner_cycle_timeout_total.load(Ordering::Relaxed);
m.cycle_recovery_required_total = self.scanner_cycle_recovery_required_total.load(Ordering::Relaxed);
m.cycle_last_progress_age = self.scanner_cycle_last_progress_age_seconds.load(Ordering::Relaxed);
m.leader_lease_without_progress = self.scanner_leader_lease_without_progress.load(Ordering::Relaxed);
m.bitrot_cycle_enabled = self.scanner_bitrot_cycle_enabled.load(Ordering::Relaxed); m.bitrot_cycle_enabled = self.scanner_bitrot_cycle_enabled.load(Ordering::Relaxed);
m.bitrot_cycle_seconds = self.scanner_bitrot_cycle_millis.load(Ordering::Relaxed) as f64 / 1000.0; m.bitrot_cycle_seconds = self.scanner_bitrot_cycle_millis.load(Ordering::Relaxed) as f64 / 1000.0;
m.scan_checkpoint = match self.scanner_checkpoint.lock() { m.scan_checkpoint = match self.scanner_checkpoint.lock() {
@@ -4966,20 +4926,4 @@ mod tests {
assert!(!report.bitrot_cycle_enabled); assert!(!report.bitrot_cycle_enabled);
assert_eq!(report.bitrot_cycle_seconds, 0.0); assert_eq!(report.bitrot_cycle_seconds, 0.0);
} }
#[tokio::test]
async fn scanner_cycle_timeout_metrics_reset_for_a_new_cycle() {
let metrics = Metrics::new();
metrics.record_scanner_cycle_timeout(true, Duration::from_secs(17));
let timed_out = metrics.report().await;
assert_eq!(timed_out.cycle_timeout_total, 1);
assert_eq!(timed_out.cycle_last_progress_age, 17);
assert!(timed_out.leader_lease_without_progress);
metrics.record_scanner_cycle_config(Duration::from_secs(60), None, Some(Duration::from_secs(1)), None, None);
let current = metrics.report().await;
assert_eq!(current.cycle_timeout_total, 1);
assert_eq!(current.cycle_last_progress_age, 0);
assert!(!current.leader_lease_without_progress);
}
} }
-6
View File
@@ -84,12 +84,6 @@ Current guidance:
- `RUSTFS_SCANNER_CYCLE_MAX_OBJECTS` (canonical) - `RUSTFS_SCANNER_CYCLE_MAX_OBJECTS` (canonical)
- `RUSTFS_SCANNER_CYCLE_MAX_DIRECTORIES` (canonical) - `RUSTFS_SCANNER_CYCLE_MAX_DIRECTORIES` (canonical)
Scanner cycle budget controls:
- When `RUSTFS_SCANNER_CYCLE_MAX_DURATION_SECS` is unset, the finite default is 1800 seconds (30 minutes), matching the scanner benchmark guidance.
- An explicit `0` preserves the compatibility behavior of an unbounded runtime budget. Object and directory budgets likewise remain unbounded when explicitly set to `0`.
- A timed-out cycle cancels cooperative scanner work, then fences its leader epoch before releasing the lease. An uncooperative I/O operation is dropped after the bounded shutdown window; its cursor is not claimed to be durable and the scanner reports `recovery-required` when the worker cannot stop cooperatively, the cycle state was not confirmed durable, or epoch fencing cannot be persisted.
## Mmap read environment aliases ## Mmap read environment aliases
- `RUSTFS_OBJECT_MMAP_READ_ENABLE` (canonical) - `RUSTFS_OBJECT_MMAP_READ_ENABLE` (canonical)
+3 -6
View File
@@ -143,12 +143,9 @@ pub const ENV_SCANNER_MAX_WAIT_SECS: &str = "RUSTFS_SCANNER_MAX_WAIT_SECS";
/// Default scanner speed preset. /// Default scanner speed preset.
pub const DEFAULT_SCANNER_SPEED: &str = "default"; pub const DEFAULT_SCANNER_SPEED: &str = "default";
/// Default scanner cycle runtime budget when no override is configured. /// Default scanner cycle runtime budget.
/// /// `0` keeps the existing unbounded per-cycle behavior.
/// An explicit `0` remains the compatibility escape hatch for an unbounded pub const DEFAULT_SCANNER_CYCLE_MAX_DURATION_SECS: u64 = 0;
/// cycle. Keeping the unset default finite prevents a stalled scanner I/O
/// operation from holding the leader lease forever.
pub const DEFAULT_SCANNER_CYCLE_MAX_DURATION_SECS: u64 = 30 * 60;
/// Default scanner per-cycle object budget. /// Default scanner per-cycle object budget.
/// `0` keeps the existing unbounded per-cycle behavior. /// `0` keeps the existing unbounded per-cycle behavior.
+2 -2
View File
@@ -317,8 +317,6 @@ pub mod config {
} }
pub mod data_usage { pub mod data_usage {
#[cfg(feature = "test-util")]
pub use crate::data_usage::seed_bucket_usage_memory_for_test;
pub use crate::data_usage::{ pub use crate::data_usage::{
DATA_USAGE_CACHE_NAME, apply_bucket_usage_memory_overlay, compute_bucket_usage, DATA_USAGE_CACHE_NAME, apply_bucket_usage_memory_overlay, compute_bucket_usage,
init_compression_total_memory_from_backend, invalidate_admin_data_usage_snapshot_cache, init_compression_total_memory_from_backend, invalidate_admin_data_usage_snapshot_cache,
@@ -330,6 +328,8 @@ pub mod data_usage {
remove_bucket_usage_from_backend, replace_bucket_usage_memory_from_info, store_compression_total_in_backend, remove_bucket_usage_from_backend, replace_bucket_usage_memory_from_info, store_compression_total_in_backend,
store_data_usage_in_backend, store_data_usage_in_backend,
}; };
#[cfg(feature = "test-util")]
pub use crate::data_usage::{get_bucket_usage_memory, seed_bucket_usage_memory_for_test};
} }
pub mod disk { pub mod disk {
@@ -2855,7 +2855,7 @@ fn replicate_object_info_from_object_info(
.map(|v| OffsetDateTime::parse(&v, &Rfc3339).unwrap_or(OffsetDateTime::UNIX_EPOCH)); .map(|v| OffsetDateTime::parse(&v, &Rfc3339).unwrap_or(OffsetDateTime::UNIX_EPOCH));
let mut rstate = oi.replication_state(); let mut rstate = oi.replication_state();
rstate.replicate_decision_str = dsc.to_string(); rstate.replicate_decision_str = dsc.to_string();
let asz = oi.get_actual_size().unwrap_or_default(); let asz = oi.get_actual_size_or_physical();
let ssec = replication_object_is_ssec_encrypted(&oi.user_defined); let ssec = replication_object_is_ssec_encrypted(&oi.user_defined);
let checksum = if ssec { oi.checksum.clone() } else { None }; let checksum = if ssec { oi.checksum.clone() } else { None };
@@ -1412,7 +1412,7 @@ pub async fn get_heal_replicate_object_info(oi: &ObjectInfo, rcfg: &ReplicationC
}; };
let mut replication_state = oi.replication_state(); let mut replication_state = oi.replication_state();
replication_state.replicate_decision_str = dsc.to_string(); replication_state.replicate_decision_str = dsc.to_string();
let actual_size = oi.get_actual_size().unwrap_or_default(); let actual_size = oi.get_actual_size_or_physical();
Ok(ReplicateObjectInfo { Ok(ReplicateObjectInfo {
name: oi.name.clone(), name: oi.name.clone(),
@@ -389,7 +389,7 @@ fn replication_source_object(object_info: &ObjectInfo) -> ReplicationSourceObjec
.map(|mod_time| OffsetDateTime::from_unix_timestamp(mod_time.unix_timestamp()).unwrap_or(mod_time)), .map(|mod_time| OffsetDateTime::from_unix_timestamp(mod_time.unix_timestamp()).unwrap_or(mod_time)),
version_id: object_info.version_id.map(|version_id| version_id.to_string()), version_id: object_info.version_id.map(|version_id| version_id.to_string()),
etag: object_info.etag.as_deref(), etag: object_info.etag.as_deref(),
actual_size: object_info.get_actual_size().unwrap_or_default(), actual_size: object_info.get_actual_size_or_physical(),
delete_marker: object_info.delete_marker, delete_marker: object_info.delete_marker,
content_type: object_info.content_type.as_deref(), content_type: object_info.content_type.as_deref(),
content_encoding: object_info.content_encoding.as_deref(), content_encoding: object_info.content_encoding.as_deref(),
@@ -542,6 +542,20 @@ mod tests {
assert!(replication_target_head_is_newer_null_version(&source, &target)); assert!(replication_target_head_is_newer_null_version(&source, &target));
} }
#[test]
fn replication_source_uses_physical_size_for_unknown_compressed_object() {
let mut metadata = HashMap::new();
rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_COMPRESSION, "zstd".to_string());
let source = ObjectInfo {
size: 128,
actual_size: -1,
user_defined: Arc::new(metadata),
..Default::default()
};
assert_eq!(replication_source_object(&source).actual_size, 128);
}
#[test] #[test]
fn replication_target_head_content_matches_compare_etag_only() { fn replication_target_head_content_matches_compare_etag_only() {
let source = ObjectInfo { let source = ObjectInfo {
+63 -60
View File
@@ -21,7 +21,7 @@ use crate::storage_api_contracts::{
bucket::{BucketInfo, BucketOperations, BucketOptions, DeleteBucketOptions, MakeBucketOptions}, bucket::{BucketInfo, BucketOperations, BucketOptions, DeleteBucketOptions, MakeBucketOptions},
list::{StorageListObjectVersionsInfo, StorageListObjectsV2Info, StorageObjectInfoOrErr, StorageWalkOptions}, list::{StorageListObjectVersionsInfo, StorageListObjectsV2Info, StorageObjectInfoOrErr, StorageWalkOptions},
multipart::{CompletePart, ListMultipartsInfo, ListPartsInfo, MultipartInfo, MultipartUploadResult, PartInfo}, multipart::{CompletePart, ListMultipartsInfo, ListPartsInfo, MultipartInfo, MultipartUploadResult, PartInfo},
object::{DeletedObject, ObjectIO as _, ObjectOperations as _, ObjectToDelete}, object::{DeleteAccounting, DeletedObject, ObjectIO as _, ObjectOperations as _, ObjectToDelete},
range::HTTPRangeSpec, range::HTTPRangeSpec,
}; };
use crate::{ use crate::{
@@ -414,6 +414,66 @@ fn apply_delete_objects_results(
} }
} }
fn apply_delete_accounting_results(
accounting: &mut [Option<DeleteAccounting>],
set_objects: &[DelObj],
set_accounting: &[Option<DeleteAccounting>],
) {
for (obj, value) in set_objects.iter().zip(set_accounting.iter()) {
accounting[obj.orig_idx] = value.clone();
}
}
impl Sets {
pub(crate) async fn delete_objects_with_accounting(
&self,
bucket: &str,
objects: Vec<ObjectToDelete>,
opts: ObjectOptions,
) -> (Vec<DeletedObject>, Vec<Option<Error>>, Vec<Option<DeleteAccounting>>) {
let mut del_objects = vec![DeletedObject::default(); objects.len()];
let mut del_errs = vec![None; objects.len()];
let mut accounting = vec![None; objects.len()];
let mut set_obj_map = HashMap::new();
for (i, obj) in objects.iter().enumerate() {
let idx = self.get_hashed_set_index(obj.object_name.as_str());
set_obj_map.entry(idx).or_insert_with(Vec::new).push(DelObj {
orig_idx: i,
obj: obj.clone(),
});
}
let max_concurrent = set_obj_map.len().min(num_cpus::get()).max(1);
let semaphore = Arc::new(tokio::sync::Semaphore::new(max_concurrent));
let mut futures = FuturesUnordered::new();
let bucket = bucket.to_owned();
for (set_index, set_objects) in set_obj_map {
let disks = self.get_disks(set_index);
let objects = set_objects.iter().map(|entry| entry.obj.clone()).collect::<Vec<_>>();
let bucket = bucket.clone();
let opts = opts.clone();
let semaphore = semaphore.clone();
futures.push(async move {
let _permit = semaphore
.acquire_owned()
.await
.expect("delete_objects semaphore should remain open");
let (deleted, errors, accounting) = disks.delete_objects_with_accounting(&bucket, objects, opts).await;
(set_objects, deleted, errors, accounting)
});
}
while let Some((set_objects, deleted, errors, set_accounting)) = futures.next().await {
apply_delete_objects_results(&mut del_objects, &mut del_errs, &set_objects, &deleted, errors);
apply_delete_accounting_results(&mut accounting, &set_objects, &set_accounting);
}
(del_objects, del_errs, accounting)
}
}
#[async_trait::async_trait] #[async_trait::async_trait]
impl crate::storage_api_contracts::object::ObjectIO for Sets { impl crate::storage_api_contracts::object::ObjectIO for Sets {
type Error = Error; type Error = Error;
@@ -655,65 +715,8 @@ impl crate::storage_api_contracts::object::ObjectOperations for Sets {
objects: Vec<ObjectToDelete>, objects: Vec<ObjectToDelete>,
opts: ObjectOptions, opts: ObjectOptions,
) -> (Vec<DeletedObject>, Vec<Option<Error>>) { ) -> (Vec<DeletedObject>, Vec<Option<Error>>) {
// Default return value let (deleted, errors, _) = self.delete_objects_with_accounting(bucket, objects, opts).await;
let mut del_objects = vec![DeletedObject::default(); objects.len()]; (deleted, errors)
let mut del_errs = Vec::with_capacity(objects.len());
for _ in 0..objects.len() {
del_errs.push(None)
}
let mut set_obj_map = HashMap::new();
// hash key
for (i, obj) in objects.iter().enumerate() {
let idx = self.get_hashed_set_index(obj.object_name.as_str());
if !set_obj_map.contains_key(&idx) {
set_obj_map.insert(
idx,
vec![DelObj {
// set_idx: idx,
orig_idx: i,
obj: obj.clone(),
}],
);
} else if let Some(val) = set_obj_map.get_mut(&idx) {
val.push(DelObj {
// set_idx: idx,
orig_idx: i,
obj: obj.clone(),
});
}
}
let max_concurrent = set_obj_map.len().min(num_cpus::get()).max(1);
let semaphore = Arc::new(tokio::sync::Semaphore::new(max_concurrent));
let mut futures = FuturesUnordered::new();
let bucket = bucket.to_string();
for (k, v) in set_obj_map {
let disks = self.get_disks(k);
let objs: Vec<ObjectToDelete> = v.iter().map(|v| v.obj.clone()).collect();
let bucket = bucket.clone();
let opts = opts.clone();
let semaphore = semaphore.clone();
futures.push(async move {
let _permit = semaphore
.acquire_owned()
.await
.expect("delete_objects semaphore should remain open");
let (dobjects, errs) = disks.delete_objects(&bucket, objs, opts).await;
(v, dobjects, errs)
});
}
while let Some((v, dobjects, errs)) = futures.next().await {
apply_delete_objects_results(&mut del_objects, &mut del_errs, &v, &dobjects, errs);
}
(del_objects, del_errs)
} }
#[tracing::instrument(skip(self))] #[tracing::instrument(skip(self))]
+108 -8
View File
@@ -1391,7 +1391,37 @@ impl BucketUsageAccumulator {
} }
pub fn quota_object_size(object: &ObjectInfo) -> Result<u64, Error> { pub fn quota_object_size(object: &ObjectInfo) -> Result<u64, Error> {
let logical_size = u64::try_from(object.get_actual_size().map_err(Error::other)?).map_err(|_| Error::PartMissingOrCorrupt)?; // A compressed object may carry -1 while the transformed size is unknown
// (legacy streaming sentinel). In that case the persisted physical size
// is still a valid accounting floor; every other negative value is corrupt.
// An explicit negative `actual-size` metadata value is corrupt, however:
// the sentinel is only valid in the in-memory/object-part field written by
// the legacy streaming path, not as a persisted declared size.
let compressed = object.is_compressed();
if object.actual_size < -1 || (object.actual_size == -1 && !compressed) {
return Err(Error::PartMissingOrCorrupt);
}
if object
.parts
.iter()
.any(|part| part.actual_size < -1 || (part.actual_size < 0 && !compressed))
{
return Err(Error::PartMissingOrCorrupt);
}
let declared_actual_size = rustfs_utils::http::get_str(&object.user_defined, rustfs_utils::http::SUFFIX_ACTUAL_SIZE)
.filter(|value| !value.is_empty());
if declared_actual_size
.as_deref()
.and_then(|value| value.parse::<i64>().ok())
.is_some_and(|size| size < 0)
{
return Err(Error::PartMissingOrCorrupt);
}
let logical_size = match object.get_actual_size().map_err(Error::other)? {
size if size == -1 && compressed && declared_actual_size.is_none() => None,
size if size >= 0 => Some(u64::try_from(size).map_err(|_| Error::PartMissingOrCorrupt)?),
_ => return Err(Error::PartMissingOrCorrupt),
};
let persisted_part_size = if object.parts.is_empty() { let persisted_part_size = if object.parts.is_empty() {
u64::try_from(object.size).map_err(|_| Error::PartMissingOrCorrupt)? u64::try_from(object.size).map_err(|_| Error::PartMissingOrCorrupt)?
} else { } else {
@@ -1399,12 +1429,8 @@ pub fn quota_object_size(object: &ObjectInfo) -> Result<u64, Error> {
// Compressed streaming objects persist -1 when the transformed // Compressed streaming objects persist -1 when the transformed
// part size is unknown. The physical part size remains a valid // part size is unknown. The physical part size remains a valid
// quota floor; reject only non-negative values that overflow. // quota floor; reject only non-negative values that overflow.
let actual_size = if part.actual_size < 0 { let actual_size = if part.actual_size == -1 {
if object.is_compressed() { 0
0
} else {
return Err(Error::PartMissingOrCorrupt);
}
} else { } else {
u64::try_from(part.actual_size).map_err(|_| Error::PartMissingOrCorrupt)? u64::try_from(part.actual_size).map_err(|_| Error::PartMissingOrCorrupt)?
}; };
@@ -1412,7 +1438,7 @@ pub fn quota_object_size(object: &ObjectInfo) -> Result<u64, Error> {
total.checked_add(part_size).ok_or(Error::PartMissingOrCorrupt) total.checked_add(part_size).ok_or(Error::PartMissingOrCorrupt)
})? })?
}; };
Ok(logical_size.max(persisted_part_size)) Ok(logical_size.unwrap_or(0).max(persisted_part_size))
} }
type UsageVersionPage = StorageListObjectVersionsInfo<ObjectInfo>; type UsageVersionPage = StorageListObjectVersionsInfo<ObjectInfo>;
@@ -3320,6 +3346,80 @@ mod tests {
); );
} }
#[test]
fn quota_object_size_accepts_compressed_unknown_actual_size_sentinel() {
let mut metadata = HashMap::new();
rustfs_utils::http::insert_str(
&mut metadata,
rustfs_utils::http::SUFFIX_COMPRESSION,
"klauspost/compress/s2".to_string(),
);
let object = ObjectInfo {
size: 400,
actual_size: -1,
user_defined: Arc::new(metadata),
..Default::default()
};
assert_eq!(quota_object_size(&object).expect("compressed sentinel is valid"), 400);
}
#[test]
fn quota_object_size_rejects_compressed_part_sum_overflow() {
let mut metadata = HashMap::new();
rustfs_utils::http::insert_str(
&mut metadata,
rustfs_utils::http::SUFFIX_COMPRESSION,
"klauspost/compress/s2".to_string(),
);
let object = ObjectInfo {
size: 1,
user_defined: Arc::new(metadata),
parts: Arc::new(vec![
rustfs_filemeta::ObjectPartInfo {
actual_size: i64::MAX,
..Default::default()
},
rustfs_filemeta::ObjectPartInfo {
actual_size: 1,
..Default::default()
},
]),
..Default::default()
};
assert!(matches!(quota_object_size(&object), Err(Error::Io(_))));
}
#[test]
fn quota_object_size_rejects_negative_values_other_than_the_compressed_sentinel() {
let mut metadata = HashMap::new();
rustfs_utils::http::insert_str(
&mut metadata,
rustfs_utils::http::SUFFIX_COMPRESSION,
"klauspost/compress/s2".to_string(),
);
let corrupt_object = ObjectInfo {
size: 400,
actual_size: -2,
user_defined: Arc::new(metadata.clone()),
..Default::default()
};
assert!(matches!(quota_object_size(&corrupt_object), Err(Error::PartMissingOrCorrupt)));
let corrupt_part = ObjectInfo {
size: 400,
user_defined: Arc::new(metadata),
parts: Arc::new(vec![rustfs_filemeta::ObjectPartInfo {
size: 400,
actual_size: -2,
..Default::default()
}]),
..Default::default()
};
assert!(matches!(quota_object_size(&corrupt_part), Err(Error::PartMissingOrCorrupt)));
}
#[tokio::test] #[tokio::test]
#[serial] #[serial]
async fn live_bucket_usage_refreshes_are_coalesced_only_while_in_flight() { async fn live_bucket_usage_refreshes_are_coalesced_only_while_in_flight() {
+34 -4
View File
@@ -689,6 +689,9 @@ impl ObjectInfo {
} }
pub fn get_actual_size(&self) -> std::io::Result<i64> { pub fn get_actual_size(&self) -> std::io::Result<i64> {
if self.actual_size < -1 || (self.actual_size == -1 && !self.is_compressed()) {
return Err(std::io::Error::other("invalid negative actual size"));
}
if self.actual_size > 0 { if self.actual_size > 0 {
return Ok(self.actual_size); return Ok(self.actual_size);
} }
@@ -700,10 +703,25 @@ impl ObjectInfo {
let size = size_str.parse::<i64>().map_err(|e| std::io::Error::other(e.to_string()))?; let size = size_str.parse::<i64>().map_err(|e| std::io::Error::other(e.to_string()))?;
return Ok(size); return Ok(size);
} }
let mut actual_size = 0; if self.actual_size == -1 && self.parts.is_empty() {
self.parts.iter().for_each(|part| { return Ok(-1);
actual_size += part.actual_size; }
}); let mut actual_size = 0_i64;
let mut unknown = false;
for part in self.parts.iter() {
match part.actual_size {
-1 => unknown = true,
size if size >= 0 => {
actual_size = actual_size
.checked_add(size)
.ok_or_else(|| std::io::Error::other("compressed actual size overflow"))?;
}
_ => return Err(std::io::Error::other("invalid negative compressed part size")),
}
}
if unknown {
return Ok(-1);
}
if actual_size == 0 && actual_size != self.size { if actual_size == 0 && actual_size != self.size {
return Err(std::io::Error::other(format!("invalid decompressed size {} {}", actual_size, self.size))); return Err(std::io::Error::other(format!("invalid decompressed size {} {}", actual_size, self.size)));
} }
@@ -718,6 +736,18 @@ impl ObjectInfo {
Ok(self.size) Ok(self.size)
} }
/// Returns a non-negative size for client and replication boundaries.
///
/// Compressed legacy metadata can retain the internal `-1` unknown-size
/// sentinel. Those boundaries cannot emit a negative length, so they use
/// the persisted physical size while quota accounting keeps the sentinel
/// distinction in [`crate::data_usage::quota_object_size`].
pub fn get_actual_size_or_physical(&self) -> i64 {
self.get_actual_size()
.map(|size| if size >= 0 { size } else { self.size.max(0) })
.unwrap_or_else(|_| self.size.max(0))
}
pub fn from_file_info(fi: &FileInfo, bucket: &str, object: &str, versioned: bool) -> ObjectInfo { pub fn from_file_info(fi: &FileInfo, bucket: &str, object: &str, versioned: bool) -> ObjectInfo {
let mut version_id = fi.version_id; let mut version_id = fi.version_id;
@@ -256,10 +256,6 @@ fn to_madmin_scanner_metrics(metrics: rustfs_common::metrics::ScannerMetricsRepo
cycle_max_duration_seconds: metrics.cycle_max_duration_seconds, cycle_max_duration_seconds: metrics.cycle_max_duration_seconds,
cycle_max_objects: metrics.cycle_max_objects, cycle_max_objects: metrics.cycle_max_objects,
cycle_max_directories: metrics.cycle_max_directories, cycle_max_directories: metrics.cycle_max_directories,
cycle_timeout_total: metrics.cycle_timeout_total,
cycle_recovery_required_total: metrics.cycle_recovery_required_total,
cycle_last_progress_age: metrics.cycle_last_progress_age,
leader_lease_without_progress: metrics.leader_lease_without_progress,
bitrot_cycle_enabled: metrics.bitrot_cycle_enabled, bitrot_cycle_enabled: metrics.bitrot_cycle_enabled,
bitrot_cycle_seconds: metrics.bitrot_cycle_seconds, bitrot_cycle_seconds: metrics.bitrot_cycle_seconds,
scan_checkpoint: metrics.scan_checkpoint.map(|checkpoint| MadminScannerCheckpointReport { scan_checkpoint: metrics.scan_checkpoint.map(|checkpoint| MadminScannerCheckpointReport {
@@ -615,10 +611,6 @@ mod test {
current_started: chrono_to_jiff_timestamp(current_started), current_started: chrono_to_jiff_timestamp(current_started),
last_cycle_partial_source: "usage".to_string(), last_cycle_partial_source: "usage".to_string(),
last_cycle_partial_source_code: 1, last_cycle_partial_source_code: 1,
cycle_timeout_total: 3,
cycle_recovery_required_total: 2,
cycle_last_progress_age: 17,
leader_lease_without_progress: true,
partial_cycles_by_source: vec![rustfs_common::metrics::ScannerSourceCycleSnapshot { partial_cycles_by_source: vec![rustfs_common::metrics::ScannerSourceCycleSnapshot {
source: "usage".to_string(), source: "usage".to_string(),
cycles: 2, cycles: 2,
@@ -630,10 +622,6 @@ mod test {
assert_eq!(scanner.current_started, chrono_to_jiff_timestamp(current_started)); assert_eq!(scanner.current_started, chrono_to_jiff_timestamp(current_started));
assert_eq!(scanner.last_cycle_partial_source, "usage"); assert_eq!(scanner.last_cycle_partial_source, "usage");
assert_eq!(scanner.last_cycle_partial_source_code, 1); assert_eq!(scanner.last_cycle_partial_source_code, 1);
assert_eq!(scanner.cycle_timeout_total, 3);
assert_eq!(scanner.cycle_recovery_required_total, 2);
assert_eq!(scanner.cycle_last_progress_age, 17);
assert!(scanner.leader_lease_without_progress);
let usage = scanner let usage = scanner
.partial_cycles_by_source .partial_cycles_by_source
.iter() .iter()
+1 -1
View File
@@ -97,7 +97,7 @@ use crate::storage_api_contracts::{
CompletePart, ListMultipartsInfo, ListPartsInfo, MultipartInfo, MultipartOperations as _, MultipartUploadResult, PartInfo, CompletePart, ListMultipartsInfo, ListPartsInfo, MultipartInfo, MultipartOperations as _, MultipartUploadResult, PartInfo,
}, },
namespace::NamespaceLocking as _, namespace::NamespaceLocking as _,
object::{DeletedObject, HTTPPreconditions, ObjectIO as _, ObjectOperations as _, ObjectToDelete}, object::{DeleteAccounting, DeletedObject, HTTPPreconditions, ObjectIO as _, ObjectOperations as _, ObjectToDelete},
range::HTTPRangeSpec, range::HTTPRangeSpec,
}; };
use crate::store::utils::is_reserved_or_invalid_bucket; use crate::store::utils::is_reserved_or_invalid_bucket;
+159 -4
View File
@@ -45,6 +45,7 @@ use crate::bucket::replication::{
DeleteReplicationConfigSnapshot, ReplicationLifecycleBridge, ReplicationStatusType, VersionPurgeStatusType, DeleteReplicationConfigSnapshot, ReplicationLifecycleBridge, ReplicationStatusType, VersionPurgeStatusType,
replication_state_to_filemeta, replication_status_from_filemeta, version_purge_status_to_filemeta, replication_state_to_filemeta, replication_status_from_filemeta, version_purge_status_to_filemeta,
}; };
use crate::data_usage::quota_object_size;
use crate::diagnostics::get::GetObjectFailureReason; use crate::diagnostics::get::GetObjectFailureReason;
use crate::disk::{DataDirDeleteStatus, OldCurrentSize}; use crate::disk::{DataDirDeleteStatus, OldCurrentSize};
use crate::error::is_err_invalid_upload_id; use crate::error::is_err_invalid_upload_id;
@@ -5655,7 +5656,18 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
objects: Vec<ObjectToDelete>, objects: Vec<ObjectToDelete>,
opts: ObjectOptions, opts: ObjectOptions,
) -> (Vec<DeletedObject>, Vec<Option<Error>>) { ) -> (Vec<DeletedObject>, Vec<Option<Error>>) {
let (deleted, errors, _) = self.delete_objects_with_accounting(bucket, objects, opts).await;
(deleted, errors)
}
async fn delete_objects_with_accounting(
&self,
bucket: &str,
objects: Vec<ObjectToDelete>,
opts: ObjectOptions,
) -> (Vec<DeletedObject>, Vec<Option<Error>>, Vec<Option<DeleteAccounting>>) {
let mut del_objects = vec![DeletedObject::default(); objects.len()]; let mut del_objects = vec![DeletedObject::default(); objects.len()];
let mut accounting = vec![None; objects.len()];
let delete_config_snapshot = opts let delete_config_snapshot = opts
.delete_replication_config_snapshot .delete_replication_config_snapshot
.clone() .clone()
@@ -5745,7 +5757,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
*item = Some(Error::other(message.clone())); *item = Some(Error::other(message.clone()));
} }
} }
return (del_objects, del_errs); return (del_objects, del_errs, accounting);
} }
}, },
} }
@@ -5792,6 +5804,22 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
let source_missing = gerr let source_missing = gerr
.as_ref() .as_ref()
.is_some_and(|err| is_err_object_not_found(err) || is_err_version_not_found(err)); .is_some_and(|err| is_err_object_not_found(err) || is_err_version_not_found(err));
// Resolve accounting from the generation selected under this
// object's write lock. A request-layer pre-stat is only an
// optimization and cannot identify a concurrent overwrite.
let (accounting_size, accounting_version_id, removed_current_object) = if source_missing
|| dobj.synthetic_version_id
|| set_disk_delete_creates_delete_marker(&check_opts)
|| goi.delete_marker
{
(None, None, false)
} else {
(
quota_object_size(&goi).ok(),
goi.version_id.filter(|version_id| !version_id.is_nil()),
(dobj.version_id.is_none() || is_explicit_null_version(dobj.version_id)) && !dobj.synthetic_version_id,
)
};
// Normalize both sides before comparing. `goi.version_id` is the // Normalize both sides before comparing. `goi.version_id` is the
// client-facing identity, where `from_file_info` synthesizes // client-facing identity, where `from_file_info` synthesizes
// `Some(Uuid::nil())` for a null version on a versioned or // `Some(Uuid::nil())` for a null version on a versioned or
@@ -5920,7 +5948,12 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
}, },
replication_state: vr.replication_state_internal.clone(), replication_state: vr.replication_state_internal.clone(),
..Default::default() ..Default::default()
} };
accounting[i] = Some(DeleteAccounting {
size: accounting_size,
version_id: accounting_version_id,
removed_current_object,
});
} }
// Only add to vers_map if we hold the lock // Only add to vers_map if we hold the lock
@@ -5966,7 +5999,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
}); });
} }
} }
return (del_objects, del_errs); return (del_objects, del_errs, accounting);
} }
let mut persisted_journal_entries = Vec::with_capacity(journal_entries.len()); let mut persisted_journal_entries = Vec::with_capacity(journal_entries.len());
@@ -6204,7 +6237,16 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
} }
} }
(del_objects, del_errs) // An accounting identity is actionable only when the delete result is
// successful. Never let a failed commit (including a partial quorum
// failure) reach the request-layer fast delta path.
for (index, err) in del_errs.iter().enumerate() {
if err.is_some() {
accounting[index] = None;
}
}
(del_objects, del_errs, accounting)
} }
#[tracing::instrument(skip(self))] #[tracing::instrument(skip(self))]
@@ -6533,6 +6575,12 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
let mut obj_info = ObjectInfo::from_file_info(&dfi, bucket, object, opts.versioned || opts.version_suspended); let mut obj_info = ObjectInfo::from_file_info(&dfi, bucket, object, opts.versioned || opts.version_suspended);
obj_info.size = goi.size; obj_info.size = goi.size;
// Keep the committed source metadata on the internal delete result so
// the request layer can derive canonical accounting for this exact
// generation. Delete responses do not expose these fields.
obj_info.actual_size = goi.actual_size;
obj_info.user_defined = Arc::clone(&goi.user_defined);
obj_info.parts = Arc::clone(&goi.parts);
obj_info.user_tags = Arc::clone(&goi.user_tags); obj_info.user_tags = Arc::clone(&goi.user_tags);
self.invalidate_get_object_metadata_cache(bucket, object).await; self.invalidate_get_object_metadata_cache(bucket, object).await;
Ok(obj_info) Ok(obj_info)
@@ -7824,6 +7872,113 @@ mod replication_quota_safety_tests {
assert_eq!(stored.get_actual_size().expect("stored logical size should parse"), 1); assert_eq!(stored.get_actual_size().expect("stored logical size should parse"), 1);
} }
#[tokio::test]
async fn delete_returns_canonical_compressed_accounting_size() {
let (_temp_dirs, disks, set_disks) = hermetic_set_disks(4).await;
let bucket = "compressed-delete-accounting";
for disk in &disks {
disk.make_volume(bucket).await.expect("bucket volume should be created");
}
let mut user_defined = HashMap::new();
insert_str(
&mut user_defined,
rustfs_utils::http::SUFFIX_COMPRESSION,
"klauspost/compress/s2".to_string(),
);
insert_str(&mut user_defined, SUFFIX_ACTUAL_SIZE, "1000".to_string());
let mut reader = PutObjReader::new(
HashReader::from_stream(Cursor::new(vec![0x5a; 400]), 400, 1000, None, None, false)
.expect("compressed fixture reader should be valid"),
);
set_disks
.put_object(
bucket,
"object",
&mut reader,
&ObjectOptions {
user_defined,
..Default::default()
},
)
.await
.expect("compressed object should be written");
let (deleted, errors, accounting) = set_disks
.delete_objects_with_accounting(
bucket,
vec![ObjectToDelete {
object_name: "object".to_string(),
..Default::default()
}],
ObjectOptions {
object_lock_config_snapshot: Some(Arc::new(ObjectLockConfigSnapshot::new(
ObjectLockConfigState::ConfirmedAbsent,
))),
..Default::default()
},
)
.await;
assert!(errors[0].is_none(), "compressed delete should succeed: {:?}", errors[0]);
assert!(deleted[0].found, "the committed object must be reported as found");
assert_eq!(accounting[0].as_ref().and_then(|value| value.size), Some(1000));
assert!(accounting[0].as_ref().is_some_and(|value| value.version_id.is_none()));
assert!(accounting[0].as_ref().is_some_and(|value| value.removed_current_object));
}
#[tokio::test]
async fn suspended_delete_marker_does_not_return_body_accounting() {
let (_temp_dirs, disks, set_disks) = hermetic_set_disks(4).await;
let bucket = "suspended-delete-accounting";
for disk in &disks {
disk.make_volume(bucket).await.expect("bucket volume should be created");
}
let mut user_defined = HashMap::new();
insert_str(
&mut user_defined,
rustfs_utils::http::SUFFIX_COMPRESSION,
"klauspost/compress/s2".to_string(),
);
insert_str(&mut user_defined, SUFFIX_ACTUAL_SIZE, "1000".to_string());
let mut reader = PutObjReader::new(
HashReader::from_stream(Cursor::new(vec![0x5a; 400]), 400, 1000, None, None, false)
.expect("compressed fixture reader should be valid"),
);
let suspended_opts = ObjectOptions {
version_suspended: true,
delete_replication_config_snapshot: Some(Arc::new(DeleteReplicationConfigSnapshot::from_configs_for_test(
s3s::dto::VersioningConfiguration {
status: Some(s3s::dto::BucketVersioningStatus::from_static(s3s::dto::BucketVersioningStatus::SUSPENDED)),
..Default::default()
},
None,
))),
user_defined,
object_lock_config_snapshot: Some(Arc::new(ObjectLockConfigSnapshot::new(ObjectLockConfigState::ConfirmedAbsent))),
..Default::default()
};
set_disks
.put_object(bucket, "object", &mut reader, &suspended_opts)
.await
.expect("compressed object should be written");
let (deleted, errors, accounting) = set_disks
.delete_objects_with_accounting(
bucket,
vec![ObjectToDelete {
object_name: "object".to_string(),
..Default::default()
}],
suspended_opts,
)
.await;
assert!(errors[0].is_none(), "suspended delete should create a marker: {:?}", errors[0]);
assert!(deleted[0].delete_marker);
assert!(accounting[0].is_none(), "a delete marker must not carry body accounting");
}
#[tokio::test] #[tokio::test]
async fn direct_put_cannot_persist_a_tiny_logical_size() { async fn direct_put_cannot_persist_a_tiny_logical_size() {
let (_temp_dirs, disks, set_disks) = hermetic_set_disks(4).await; let (_temp_dirs, disks, set_disks) = hermetic_set_disks(4).await;
@@ -62,8 +62,8 @@ pub(crate) mod object {
use super::{Debug, Error, FileInfo, GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader}; use super::{Debug, Error, FileInfo, GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader};
use crate::storage_api_contracts::range::HTTPRangeSpec; use crate::storage_api_contracts::range::HTTPRangeSpec;
pub(crate) use rustfs_storage_api::{ pub(crate) use rustfs_storage_api::{
DeletedObject, HTTPPreconditions, ObjectIO, ObjectLockDeleteOptions, ObjectLockRetentionOptions, ObjectOperations, DeleteAccounting, DeletedObject, HTTPPreconditions, ObjectIO, ObjectLockDeleteOptions, ObjectLockRetentionOptions,
ObjectPreconditionError, ObjectPreconditionPart, ObjectPreconditionState, ObjectToDelete, ObjectOperations, ObjectPreconditionError, ObjectPreconditionPart, ObjectPreconditionState, ObjectToDelete,
}; };
pub(crate) trait EcstoreObjectIO: pub(crate) trait EcstoreObjectIO:
+54 -11
View File
@@ -41,7 +41,7 @@ use crate::set_disk::{
}; };
use crate::storage_api_contracts::{ use crate::storage_api_contracts::{
namespace::NamespaceLocking as _, namespace::NamespaceLocking as _,
object::{ObjectIO as _, ObjectOperations as _}, object::{DeleteAccounting, ObjectIO as _, ObjectOperations as _},
}; };
use parking_lot::Mutex as ParkingMutex; use parking_lot::Mutex as ParkingMutex;
use rustfs_io_metrics::{ use rustfs_io_metrics::{
@@ -1216,6 +1216,14 @@ fn return_batch_delete_lock_error(objects: &[ObjectToDelete], err: Error) -> (Ve
(del_objects, del_errs) (del_objects, del_errs)
} }
fn return_batch_delete_lock_error_with_accounting(
objects: &[ObjectToDelete],
err: Error,
) -> (Vec<DeletedObject>, Vec<Option<Error>>, Vec<Option<DeleteAccounting>>) {
let (deleted, errors) = return_batch_delete_lock_error(objects, err);
(deleted, errors, vec![None; objects.len()])
}
fn sorted_unique_delete_object_names(objects: &[ObjectToDelete]) -> Vec<&str> { fn sorted_unique_delete_object_names(objects: &[ObjectToDelete]) -> Vec<&str> {
let mut object_names: Vec<&str> = objects.iter().map(|object| object.object_name.as_str()).collect(); let mut object_names: Vec<&str> = objects.iter().map(|object| object.object_name.as_str()).collect();
object_names.sort_unstable(); object_names.sort_unstable();
@@ -2312,6 +2320,22 @@ impl ECStore {
result result
} }
pub async fn delete_objects_with_tier_delete_journal_and_accounting(
self: &Arc<Self>,
bucket: &str,
objects: Vec<ObjectToDelete>,
opts: ObjectOptions,
) -> (Vec<DeletedObject>, Vec<Option<Error>>, Vec<Option<DeleteAccounting>>) {
let result = self
.handle_delete_objects_with_journal_and_accounting(bucket, objects, opts, Some(Arc::clone(self)))
.await;
let success_count = result.1.iter().filter(|err| err.is_none()).count();
if success_count > 0 {
list_objects::observe_list_objects_mutations(self, bucket, success_count).await;
}
result
}
#[instrument(skip(self))] #[instrument(skip(self))]
pub(super) async fn handle_delete_object(&self, bucket: &str, object: &str, opts: ObjectOptions) -> Result<ObjectInfo> { pub(super) async fn handle_delete_object(&self, bucket: &str, object: &str, opts: ObjectOptions) -> Result<ObjectInfo> {
self.handle_delete_object_with_journal(bucket, object, opts, None).await self.handle_delete_object_with_journal(bucket, object, opts, None).await
@@ -2689,6 +2713,19 @@ impl ECStore {
opts: ObjectOptions, opts: ObjectOptions,
tier_journal_api: Option<Arc<ECStore>>, tier_journal_api: Option<Arc<ECStore>>,
) -> (Vec<DeletedObject>, Vec<Option<Error>>) { ) -> (Vec<DeletedObject>, Vec<Option<Error>>) {
let (deleted, errors, _) = self
.handle_delete_objects_with_journal_and_accounting(bucket, objects, opts, tier_journal_api)
.await;
(deleted, errors)
}
pub(super) async fn handle_delete_objects_with_journal_and_accounting(
&self,
bucket: &str,
objects: Vec<ObjectToDelete>,
opts: ObjectOptions,
tier_journal_api: Option<Arc<ECStore>>,
) -> (Vec<DeletedObject>, Vec<Option<Error>>, Vec<Option<DeleteAccounting>>) {
// encode object name // encode object name
let objects: Vec<ObjectToDelete> = objects let objects: Vec<ObjectToDelete> = objects
.iter() .iter()
@@ -2701,6 +2738,7 @@ impl ECStore {
// Default return value // Default return value
let mut del_objects = vec![DeletedObject::default(); objects.len()]; let mut del_objects = vec![DeletedObject::default(); objects.len()];
let mut accounting = vec![None; objects.len()];
let mut del_errs = Vec::with_capacity(objects.len()); let mut del_errs = Vec::with_capacity(objects.len());
for _ in 0..objects.len() { for _ in 0..objects.len() {
@@ -2714,7 +2752,7 @@ impl ECStore {
} else { } else {
match self.acquire_bucket_lifecycle_read_lock(bucket).await { match self.acquire_bucket_lifecycle_read_lock(bucket).await {
Ok(guard) => Some(guard), Ok(guard) => Some(guard),
Err(err) => return return_batch_delete_lock_error(objects.as_slice(), err), Err(err) => return return_batch_delete_lock_error_with_accounting(objects.as_slice(), err),
} }
}; };
if let Some(guard) = _bucket_lifecycle_guard.as_ref() { if let Some(guard) = _bucket_lifecycle_guard.as_ref() {
@@ -2726,21 +2764,21 @@ impl ECStore {
Err(err) => { Err(err) => {
let message = err.to_string(); let message = err.to_string();
let errors = (0..objects.len()).map(|_| Some(Error::other(message.clone()))).collect(); let errors = (0..objects.len()).map(|_| Some(Error::other(message.clone()))).collect();
return (del_objects, errors); return (del_objects, errors, accounting);
} }
} }
} }
if !is_meta_bucketname(bucket) if !is_meta_bucketname(bucket)
&& let Err(err) = get_cached_bucket_incarnation_id_in(&self.ctx, bucket).await && let Err(err) = get_cached_bucket_incarnation_id_in(&self.ctx, bucket).await
{ {
return return_batch_delete_lock_error(objects.as_slice(), err); return return_batch_delete_lock_error_with_accounting(objects.as_slice(), err);
} }
let _object_lock_metadata_guard = if is_meta_bucketname(bucket) { let _object_lock_metadata_guard = if is_meta_bucketname(bucket) {
None None
} else { } else {
Some(match acquire_bucket_metadata_transaction_read_lock_in(&self.ctx, bucket).await { Some(match acquire_bucket_metadata_transaction_read_lock_in(&self.ctx, bucket).await {
Ok(guard) => guard, Ok(guard) => guard,
Err(err) => return return_batch_delete_lock_error(objects.as_slice(), err), Err(err) => return return_batch_delete_lock_error_with_accounting(objects.as_slice(), err),
}) })
}; };
if let Some(guard) = _object_lock_metadata_guard.as_ref() { if let Some(guard) = _object_lock_metadata_guard.as_ref() {
@@ -2750,7 +2788,7 @@ impl ECStore {
let (state, incarnation_id, config_revision) = let (state, incarnation_id, config_revision) =
match get_object_lock_config_and_incarnation_from_disk_in(&self.ctx, bucket).await { match get_object_lock_config_and_incarnation_from_disk_in(&self.ctx, bucket).await {
Ok(snapshot) => snapshot, Ok(snapshot) => snapshot,
Err(err) => return return_batch_delete_lock_error(objects.as_slice(), err), Err(err) => return return_batch_delete_lock_error_with_accounting(objects.as_slice(), err),
}; };
opts.object_lock_config_snapshot = Some(Arc::new(ObjectLockConfigSnapshot::for_store_bucket( opts.object_lock_config_snapshot = Some(Arc::new(ObjectLockConfigSnapshot::for_store_bucket(
self.id, self.id,
@@ -2766,7 +2804,10 @@ impl ECStore {
if let (Some(expected), Some(current)) = (opts.expected_bucket_incarnation_id, current_bucket_incarnation_id) if let (Some(expected), Some(current)) = (opts.expected_bucket_incarnation_id, current_bucket_incarnation_id)
&& expected != current && expected != current
{ {
return return_batch_delete_lock_error(objects.as_slice(), StorageError::BucketNotFound(bucket.to_string())); return return_batch_delete_lock_error_with_accounting(
objects.as_slice(),
StorageError::BucketNotFound(bucket.to_string()),
);
} }
#[cfg(test)] #[cfg(test)]
if current_bucket_incarnation_id.is_some() { if current_bucket_incarnation_id.is_some() {
@@ -2774,7 +2815,7 @@ impl ECStore {
} }
let _object_lock_guards = match self.acquire_delete_objects_write_locks(bucket, &objects, &mut opts).await { let _object_lock_guards = match self.acquire_delete_objects_write_locks(bucket, &objects, &mut opts).await {
Ok(guards) => guards, Ok(guards) => guards,
Err(err) => return return_batch_delete_lock_error(objects.as_slice(), err), Err(err) => return return_batch_delete_lock_error_with_accounting(objects.as_slice(), err),
}; };
let mut futures = Vec::with_capacity(self.pools.len()); let mut futures = Vec::with_capacity(self.pools.len());
@@ -2783,22 +2824,24 @@ impl ECStore {
if self.is_pool_rebalancing(pool.pool_idx).await { if self.is_pool_rebalancing(pool.pool_idx).await {
continue; continue;
} }
futures.push(pool.delete_objects(bucket, objects.clone(), opts.clone())); futures.push(pool.delete_objects_with_accounting(bucket, objects.clone(), opts.clone()));
} }
let results = join_all(futures).await; let results = join_all(futures).await;
for idx in 0..del_objects.len() { for idx in 0..del_objects.len() {
for (dels, errs) in results.iter() { for (dels, errs, pool_accounting) in results.iter() {
if errs[idx].is_none() && dels[idx].found { if errs[idx].is_none() && dels[idx].found {
del_errs[idx] = None; del_errs[idx] = None;
del_objects[idx] = dels[idx].clone(); del_objects[idx] = dels[idx].clone();
accounting[idx] = pool_accounting[idx].clone();
break; break;
} }
if del_errs[idx].is_none() { if del_errs[idx].is_none() {
del_errs[idx] = errs[idx].clone(); del_errs[idx] = errs[idx].clone();
del_objects[idx] = dels[idx].clone(); del_objects[idx] = dels[idx].clone();
accounting[idx] = pool_accounting[idx].clone();
} }
} }
} }
@@ -2807,7 +2850,7 @@ impl ECStore {
v.object_name = decode_dir_object(&v.object_name); v.object_name = decode_dir_object(&v.object_name);
}); });
(del_objects, del_errs) (del_objects, del_errs, accounting)
// let mut futures = Vec::with_capacity(objects.len()); // let mut futures = Vec::with_capacity(objects.len());
-16
View File
@@ -689,14 +689,6 @@ pub struct ScannerMetrics {
pub cycle_max_objects: u64, pub cycle_max_objects: u64,
#[serde(rename = "cycle_max_directories", default)] #[serde(rename = "cycle_max_directories", default)]
pub cycle_max_directories: u64, pub cycle_max_directories: u64,
#[serde(rename = "cycle_timeout_total", default)]
pub cycle_timeout_total: u64,
#[serde(rename = "cycle_recovery_required_total", default)]
pub cycle_recovery_required_total: u64,
#[serde(rename = "cycle_last_progress_age", default)]
pub cycle_last_progress_age: u64,
#[serde(rename = "leader_lease_without_progress", default)]
pub leader_lease_without_progress: bool,
#[serde(rename = "bitrot_cycle_enabled", default)] #[serde(rename = "bitrot_cycle_enabled", default)]
pub bitrot_cycle_enabled: bool, pub bitrot_cycle_enabled: bool,
#[serde(rename = "bitrot_cycle_seconds", default)] #[serde(rename = "bitrot_cycle_seconds", default)]
@@ -772,8 +764,6 @@ impl ScannerMetrics {
self.cycle_max_duration_seconds = other.cycle_max_duration_seconds; self.cycle_max_duration_seconds = other.cycle_max_duration_seconds;
self.cycle_max_objects = other.cycle_max_objects; self.cycle_max_objects = other.cycle_max_objects;
self.cycle_max_directories = other.cycle_max_directories; self.cycle_max_directories = other.cycle_max_directories;
self.cycle_last_progress_age = other.cycle_last_progress_age;
self.leader_lease_without_progress = other.leader_lease_without_progress;
self.bitrot_cycle_enabled = other.bitrot_cycle_enabled; self.bitrot_cycle_enabled = other.bitrot_cycle_enabled;
self.bitrot_cycle_seconds = other.bitrot_cycle_seconds; self.bitrot_cycle_seconds = other.bitrot_cycle_seconds;
} }
@@ -867,12 +857,6 @@ impl ScannerMetrics {
.saturating_add(other.last_cycle_replication_checks); .saturating_add(other.last_cycle_replication_checks);
self.last_cycle_usage_saves = self.last_cycle_usage_saves.saturating_add(other.last_cycle_usage_saves); self.last_cycle_usage_saves = self.last_cycle_usage_saves.saturating_add(other.last_cycle_usage_saves);
self.failed_cycles = self.failed_cycles.saturating_add(other.failed_cycles); self.failed_cycles = self.failed_cycles.saturating_add(other.failed_cycles);
self.cycle_timeout_total = self.cycle_timeout_total.saturating_add(other.cycle_timeout_total);
self.cycle_recovery_required_total = self
.cycle_recovery_required_total
.saturating_add(other.cycle_recovery_required_total);
self.cycle_last_progress_age = self.cycle_last_progress_age.max(other.cycle_last_progress_age);
self.leader_lease_without_progress |= other.leader_lease_without_progress;
self.superseded_cycles = self.superseded_cycles.saturating_add(other.superseded_cycles); self.superseded_cycles = self.superseded_cycles.saturating_add(other.superseded_cycles);
self.partial_cycles_unknown = self.partial_cycles_unknown.saturating_add(other.partial_cycles_unknown); self.partial_cycles_unknown = self.partial_cycles_unknown.saturating_add(other.partial_cycles_unknown);
self.partial_cycles_runtime = self.partial_cycles_runtime.saturating_add(other.partial_cycles_runtime); self.partial_cycles_runtime = self.partial_cycles_runtime.saturating_add(other.partial_cycles_runtime);
+23 -95
View File
@@ -125,10 +125,7 @@ impl Default for ScannerRuntimeConfig {
cycle_interval_source: ScannerRuntimeConfigSource::Default, cycle_interval_source: ScannerRuntimeConfigSource::Default,
bitrot_cycle: Some(Duration::from_secs(DEFAULT_HEAL_BITROT_CYCLE_SECS)), bitrot_cycle: Some(Duration::from_secs(DEFAULT_HEAL_BITROT_CYCLE_SECS)),
bitrot_cycle_source: ScannerRuntimeConfigSource::Default, bitrot_cycle_source: ScannerRuntimeConfigSource::Default,
cycle_budget: ScannerCycleBudgetConfig { cycle_budget: ScannerCycleBudgetConfig::default(),
max_duration: Some(Duration::from_secs(DEFAULT_SCANNER_CYCLE_MAX_DURATION_SECS)),
..Default::default()
},
cycle_max_duration_source: ScannerRuntimeConfigSource::Default, cycle_max_duration_source: ScannerRuntimeConfigSource::Default,
cycle_max_objects_source: ScannerRuntimeConfigSource::Default, cycle_max_objects_source: ScannerRuntimeConfigSource::Default,
cycle_max_directories_source: ScannerRuntimeConfigSource::Default, cycle_max_directories_source: ScannerRuntimeConfigSource::Default,
@@ -377,10 +374,7 @@ fn validate_persisted_scanner_runtime_config(config: &ServerConfig) -> Result<()
} }
validate_optional_config_u64(scanner_kvs, SCANNER_START_DELAY, "")?; validate_optional_config_u64(scanner_kvs, SCANNER_START_DELAY, "")?;
validate_optional_config_u64(scanner_kvs, SCANNER_CYCLE, "")?; validate_optional_config_u64(scanner_kvs, SCANNER_CYCLE, "")?;
if let Some(value) = config_value(scanner_kvs, SCANNER_CYCLE_MAX_DURATION, DEFAULT_SCANNER_CYCLE_MAX_DURATION_SECS) { validate_optional_config_u64(scanner_kvs, SCANNER_CYCLE_MAX_DURATION, DEFAULT_SCANNER_CYCLE_MAX_DURATION_SECS)?;
let secs = parse_config_u64(SCANNER_CYCLE_MAX_DURATION, value)?;
cycle_duration_from_secs(SCANNER_CYCLE_MAX_DURATION, secs)?;
}
validate_optional_config_u64(scanner_kvs, SCANNER_CYCLE_MAX_OBJECTS, DEFAULT_SCANNER_CYCLE_MAX_OBJECTS)?; validate_optional_config_u64(scanner_kvs, SCANNER_CYCLE_MAX_OBJECTS, DEFAULT_SCANNER_CYCLE_MAX_OBJECTS)?;
validate_optional_config_u64(scanner_kvs, SCANNER_CYCLE_MAX_DIRECTORIES, DEFAULT_SCANNER_CYCLE_MAX_DIRECTORIES)?; validate_optional_config_u64(scanner_kvs, SCANNER_CYCLE_MAX_DIRECTORIES, DEFAULT_SCANNER_CYCLE_MAX_DIRECTORIES)?;
if let Some(value) = config_value(heal_kvs, HEAL_BITROT_CYCLE, DEFAULT_HEAL_BITROT_CYCLE_SECS) { if let Some(value) = config_value(heal_kvs, HEAL_BITROT_CYCLE, DEFAULT_HEAL_BITROT_CYCLE_SECS) {
@@ -442,46 +436,19 @@ fn lookup_max_wait(
Ok((speed.max_sleep(), speed_source)) Ok((speed.max_sleep(), speed_source))
} }
fn lookup_cycle_duration(kvs: Option<&KVS>) -> Result<(Option<Duration>, ScannerRuntimeConfigSource), ScannerRuntimeConfigError> { fn lookup_optional_seconds(
match rustfs_utils::get_env_parse_outcome::<u64>(ENV_SCANNER_CYCLE_MAX_DURATION_SECS) { kvs: Option<&KVS>,
rustfs_utils::EnvParseOutcome::Parsed(secs) => { key: &'static str,
return cycle_duration_from_secs(ENV_SCANNER_CYCLE_MAX_DURATION_SECS, secs) env_key: &'static str,
.map(|duration| (duration, ScannerRuntimeConfigSource::Env)); default: u64,
} ) -> Result<(Option<Duration>, ScannerRuntimeConfigSource), ScannerRuntimeConfigError> {
rustfs_utils::EnvParseOutcome::Invalid => { if let Some(secs) = rustfs_utils::get_env_opt_u64(env_key) {
// Do not include the raw environment value in the typed error: return Ok((Some(Duration::from_secs(secs)), ScannerRuntimeConfigSource::Env));
// deployments occasionally put sensitive material in inherited
// environment snapshots. The key still identifies the control.
return Err(invalid_value(
ENV_SCANNER_CYCLE_MAX_DURATION_SECS,
"<invalid>",
"expected unsigned integer seconds",
));
}
rustfs_utils::EnvParseOutcome::Absent => {}
} }
if let Some(value) = config_value(kvs, key, default) {
if let Some(value) = config_value(kvs, SCANNER_CYCLE_MAX_DURATION, DEFAULT_SCANNER_CYCLE_MAX_DURATION_SECS) { return parse_config_u64(key, value).map(|secs| (Some(Duration::from_secs(secs)), ScannerRuntimeConfigSource::Config));
let secs = parse_config_u64(SCANNER_CYCLE_MAX_DURATION, value)?;
return cycle_duration_from_secs(SCANNER_CYCLE_MAX_DURATION, secs)
.map(|duration| (duration, ScannerRuntimeConfigSource::Config));
} }
Ok((None, ScannerRuntimeConfigSource::Default))
Ok((
Some(Duration::from_secs(DEFAULT_SCANNER_CYCLE_MAX_DURATION_SECS)),
ScannerRuntimeConfigSource::Default,
))
}
fn cycle_duration_from_secs(key: &'static str, secs: u64) -> Result<Option<Duration>, ScannerRuntimeConfigError> {
if secs == 0 {
return Ok(None);
}
let duration = Duration::from_secs(secs);
if std::time::Instant::now().checked_add(duration).is_none() {
return Err(invalid_value(key, "<overflow>", "duration exceeds the timer range"));
}
Ok(Some(duration))
} }
fn lookup_start_delay(kvs: Option<&KVS>) -> Result<(Option<Duration>, ScannerRuntimeConfigSource), ScannerRuntimeConfigError> { fn lookup_start_delay(kvs: Option<&KVS>) -> Result<(Option<Duration>, ScannerRuntimeConfigSource), ScannerRuntimeConfigError> {
@@ -586,7 +553,12 @@ pub(crate) fn lookup_scanner_runtime_config(
(speed.cycle_interval(), speed_source) (speed.cycle_interval(), speed_source)
}; };
let (cycle_max_duration, cycle_max_duration_source) = lookup_cycle_duration(scanner_kvs)?; let (cycle_max_duration, cycle_max_duration_source) = lookup_optional_seconds(
scanner_kvs,
SCANNER_CYCLE_MAX_DURATION,
ENV_SCANNER_CYCLE_MAX_DURATION_SECS,
DEFAULT_SCANNER_CYCLE_MAX_DURATION_SECS,
)?;
let (cycle_max_objects, cycle_max_objects_source) = lookup_count_budget( let (cycle_max_objects, cycle_max_objects_source) = lookup_count_budget(
scanner_kvs, scanner_kvs,
SCANNER_CYCLE_MAX_OBJECTS, SCANNER_CYCLE_MAX_OBJECTS,
@@ -891,10 +863,10 @@ mod tests {
use rustfs_config::server_config::{Config as ServerConfig, KVS}; use rustfs_config::server_config::{Config as ServerConfig, KVS};
use rustfs_config::{ use rustfs_config::{
DEFAULT_DELIMITER, DEFAULT_HEAL_BITROT_CYCLE_SECS, ENV_SCANNER_BITROT_CYCLE_SECS, ENV_SCANNER_CACHE_SAVE_TIMEOUT_SECS, DEFAULT_DELIMITER, DEFAULT_HEAL_BITROT_CYCLE_SECS, ENV_SCANNER_BITROT_CYCLE_SECS, ENV_SCANNER_CACHE_SAVE_TIMEOUT_SECS,
ENV_SCANNER_CYCLE, ENV_SCANNER_CYCLE_MAX_DURATION_SECS, ENV_SCANNER_CYCLE_MAX_OBJECTS, ENV_SCANNER_DELAY, ENV_SCANNER_CYCLE, ENV_SCANNER_CYCLE_MAX_OBJECTS, ENV_SCANNER_DELAY, ENV_SCANNER_MAX_WAIT_SECS, ENV_SCANNER_SPEED,
ENV_SCANNER_MAX_WAIT_SECS, ENV_SCANNER_SPEED, HEAL_BITROT_CYCLE, HEAL_SUB_SYS, SCANNER_BITROT_CYCLE, HEAL_BITROT_CYCLE, HEAL_SUB_SYS, SCANNER_BITROT_CYCLE, SCANNER_CACHE_SAVE_TIMEOUT, SCANNER_CYCLE,
SCANNER_CACHE_SAVE_TIMEOUT, SCANNER_CYCLE, SCANNER_CYCLE_MAX_DIRECTORIES, SCANNER_CYCLE_MAX_DURATION, SCANNER_CYCLE_MAX_DIRECTORIES, SCANNER_CYCLE_MAX_DURATION, SCANNER_CYCLE_MAX_OBJECTS, SCANNER_DELAY, SCANNER_IDLE_MODE,
SCANNER_CYCLE_MAX_OBJECTS, SCANNER_DELAY, SCANNER_IDLE_MODE, SCANNER_SPEED, SCANNER_SUB_SYS, ScannerSpeed, SCANNER_SPEED, SCANNER_SUB_SYS, ScannerSpeed,
}; };
use std::collections::HashMap; use std::collections::HashMap;
use std::time::Duration; use std::time::Duration;
@@ -969,50 +941,6 @@ mod tests {
}); });
} }
#[test]
fn scanner_unset_budget_uses_safe_default_but_explicit_zero_is_unbounded() {
let config = server_config_with_scanner(&[]);
with_var_unset(ENV_SCANNER_CYCLE_MAX_DURATION_SECS, || {
let resolved = lookup_scanner_runtime_config(Some(&config)).expect("scanner runtime config");
assert_eq!(resolved.cycle_budget.max_duration, Some(Duration::from_secs(1800)));
assert_eq!(resolved.cycle_max_duration_source, ScannerRuntimeConfigSource::Default);
});
let config = server_config_with_scanner(&[(SCANNER_CYCLE_MAX_DURATION, "0")]);
with_var_unset(ENV_SCANNER_CYCLE_MAX_DURATION_SECS, || {
let resolved = lookup_scanner_runtime_config(Some(&config)).expect("scanner runtime config");
assert_eq!(resolved.cycle_budget.max_duration, None);
assert_eq!(resolved.cycle_max_duration_source, ScannerRuntimeConfigSource::Config);
});
}
#[test]
fn cycle_budget_invalid_or_overflow_config_is_rejected() {
with_var(ENV_SCANNER_CYCLE_MAX_DURATION_SECS, Some("invalid"), || {
let error = lookup_scanner_runtime_config(None).expect_err("invalid duration env must be rejected");
assert!(error.to_string().contains(ENV_SCANNER_CYCLE_MAX_DURATION_SECS));
assert!(error.to_string().contains("<invalid>"));
assert!(!error.to_string().contains(": invalid ("));
});
with_var(ENV_SCANNER_CYCLE_MAX_DURATION_SECS, Some("18446744073709551616"), || {
assert!(lookup_scanner_runtime_config(None).is_err());
});
with_var(ENV_SCANNER_CYCLE_MAX_DURATION_SECS, Some("18446744073709551615"), || {
assert!(lookup_scanner_runtime_config(None).is_err());
});
let config = server_config_with_scanner(&[(SCANNER_CYCLE_MAX_DURATION, "not-a-duration")]);
assert!(lookup_scanner_runtime_config(Some(&config)).is_err());
}
#[test]
fn scanner_runtime_config_validation_rejects_overflow_persisted_duration() {
let config = server_config_with_scanner(&[(SCANNER_CYCLE_MAX_DURATION, "18446744073709551615")]);
let error = validate_scanner_runtime_config(&config)
.expect_err("persisted duration that exceeds the timer range must be rejected");
assert!(error.to_string().contains(SCANNER_CYCLE_MAX_DURATION));
}
#[test] #[test]
fn scanner_runtime_config_normalizes_persisted_default_speed() { fn scanner_runtime_config_normalizes_persisted_default_speed() {
let config = server_config_with_scanner(&[(SCANNER_SPEED, "default")]); let config = server_config_with_scanner(&[(SCANNER_SPEED, "default")]);
+20 -197
View File
@@ -52,7 +52,6 @@ use rustfs_config::{
}; };
use rustfs_config::{ENV_SCANNER_CYCLE, ENV_SCANNER_SPEED, ENV_SCANNER_START_DELAY_SECS}; use rustfs_config::{ENV_SCANNER_CYCLE, ENV_SCANNER_SPEED, ENV_SCANNER_START_DELAY_SECS};
use rustfs_data_usage::observed_data_usage_is_newer; use rustfs_data_usage::observed_data_usage_is_newer;
use rustfs_lock::NamespaceLockGuard;
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use sha2::{Digest as _, Sha256}; use sha2::{Digest as _, Sha256};
use tokio::sync::{Notify, mpsc}; use tokio::sync::{Notify, mpsc};
@@ -1038,116 +1037,20 @@ fn data_usage_persist_timeout() -> Duration {
DataUsageCache::persistence_timeout() DataUsageCache::persistence_timeout()
} }
#[cfg(not(test))]
const SCANNER_CYCLE_EPOCH_FENCE_TIMEOUT: Duration = Duration::from_secs(30);
#[cfg(test)]
const SCANNER_CYCLE_EPOCH_FENCE_TIMEOUT: Duration = Duration::from_millis(50);
async fn fence_scanner_epoch_after_cycle_timeout<Store, LockLost>(
ctx: &CancellationToken,
storeapi: Arc<Store>,
cycle_info: &mut CurrentCycle,
cycle_revision: &mut DataUsageCacheRevision,
leader_epoch: &mut u64,
lock_lost: LockLost,
) -> bool
where
Store: ScannerObjectIO,
LockLost: Future<Output = ()>,
{
let fence_ctx = ctx.child_token();
let claim = claim_scanner_leadership(&fence_ctx, storeapi, cycle_info, cycle_revision, leader_epoch);
tokio::pin!(claim);
tokio::pin!(lock_lost);
tokio::select! {
biased;
_ = &mut lock_lost => {
fence_ctx.cancel();
false
}
result = tokio::time::timeout(SCANNER_CYCLE_EPOCH_FENCE_TIMEOUT, &mut claim) => {
result.unwrap_or(false) && !fence_ctx.is_cancelled()
}
}
}
struct ScannerCycleDeadlineState<'a> {
cycle_info: &'a mut CurrentCycle,
cycle_revision: &'a mut DataUsageCacheRevision,
leader_epoch: &'a mut u64,
cycle_budget: &'a ScannerCycleBudget,
}
fn cycle_timeout_requires_recovery(worker_stopped: bool, cycle_state_persisted: bool, generation_fenced: bool) -> bool {
!worker_stopped || !cycle_state_persisted || !generation_fenced
}
async fn handle_scanner_cycle_deadline<Store>(
ctx: &CancellationToken,
storeapi: Arc<Store>,
state: ScannerCycleDeadlineState<'_>,
worker_stopped: bool,
guard: &mut NamespaceLockGuard,
) where
Store: ScannerObjectIO,
{
let fenced = fence_scanner_epoch_after_cycle_timeout(
ctx,
storeapi,
state.cycle_info,
state.cycle_revision,
state.leader_epoch,
guard.lock_lost_notified(),
)
.await;
let cycle_state_persisted = state.cycle_budget.cycle_state_persisted();
let recovery_required = cycle_timeout_requires_recovery(worker_stopped, cycle_state_persisted, fenced);
warn!(
target: "rustfs::scanner",
event = EVENT_SCANNER_CYCLE_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_RUNTIME,
state = "cycle_timeout",
worker_stopped,
cycle_state_persisted,
generation_fenced = fenced,
recovery_required,
"Scanner cycle deadline expired; durable cursor/generation fencing completed when possible"
);
global_metrics().record_scanner_cycle_timeout(recovery_required, state.cycle_budget.progress_age());
// Stop renewing before releasing the lease. A new leader can then claim the
// higher persisted generation instead of inheriting the expired worker.
guard.release();
global_metrics().set_cycle(None).await;
}
async fn mark_scan_cycle_idle(cycle_info: &mut CurrentCycle, cycle_metrics_guard: &mut ScannerCycleMetricsGuard) { async fn mark_scan_cycle_idle(cycle_info: &mut CurrentCycle, cycle_metrics_guard: &mut ScannerCycleMetricsGuard) {
cycle_info.current = 0; cycle_info.current = 0;
global_metrics().clear_current_scan_mode(); global_metrics().clear_current_scan_mode();
cycle_metrics_guard.finish(cycle_info.clone()).await; cycle_metrics_guard.finish(cycle_info.clone()).await;
} }
#[cfg(test)] #[instrument(skip_all)]
#[hotpath::measure]
async fn run_data_scanner_cycle( async fn run_data_scanner_cycle(
ctx: &CancellationToken, ctx: &CancellationToken,
storeapi: &Arc<ECStore>, storeapi: &Arc<ECStore>,
cycle_info: &mut CurrentCycle, cycle_info: &mut CurrentCycle,
cycle_revision: &mut DataUsageCacheRevision, cycle_revision: &mut DataUsageCacheRevision,
leader_epoch: u64, leader_epoch: u64,
) -> ScannerCycleOutcome {
let cycle_budget = ScannerCycleBudget::new(ctx, scanner_cycle_budget_config());
run_data_scanner_cycle_with_budget(ctx, storeapi, cycle_info, cycle_revision, leader_epoch, cycle_budget).await
}
#[instrument(skip_all)]
#[hotpath::measure]
async fn run_data_scanner_cycle_with_budget(
ctx: &CancellationToken,
storeapi: &Arc<ECStore>,
cycle_info: &mut CurrentCycle,
cycle_revision: &mut DataUsageCacheRevision,
leader_epoch: u64,
cycle_budget: Arc<ScannerCycleBudget>,
) -> ScannerCycleOutcome { ) -> ScannerCycleOutcome {
let _activity_guard = ScannerActivityGuard::new(); let _activity_guard = ScannerActivityGuard::new();
if let Err(err) = refresh_scanner_runtime_config_from_global() { if let Err(err) = refresh_scanner_runtime_config_from_global() {
@@ -1163,11 +1066,7 @@ async fn run_data_scanner_cycle_with_budget(
} }
let configured_cycle_interval = scanner_cycle_interval(); let configured_cycle_interval = scanner_cycle_interval();
let configured_bitrot_cycle = scanner_bitrot_cycle(); let configured_bitrot_cycle = scanner_bitrot_cycle();
let cycle_budget_config = ScannerCycleBudgetConfig { let cycle_budget_config = scanner_cycle_budget_config();
max_duration: cycle_budget.max_duration(),
max_objects: cycle_budget.max_objects(),
max_directories: cycle_budget.max_directories(),
};
let usage_persist_timeout = data_usage_persist_timeout(); let usage_persist_timeout = data_usage_persist_timeout();
global_metrics().record_scanner_cycle_config( global_metrics().record_scanner_cycle_config(
configured_cycle_interval, configured_cycle_interval,
@@ -1238,6 +1137,7 @@ async fn run_data_scanner_cycle_with_budget(
let (sender, receiver) = mpsc::channel::<DataUsageInfo>(1); let (sender, receiver) = mpsc::channel::<DataUsageInfo>(1);
let done_cycle = Metrics::time(Metric::ScanCycle); let done_cycle = Metrics::time(Metric::ScanCycle);
let cycle_budget = ScannerCycleBudget::new(ctx, cycle_budget_config);
let scan_result = storeapi let scan_result = storeapi
.clone() .clone()
.nsscanner_with_status( .nsscanner_with_status(
@@ -1377,7 +1277,7 @@ async fn run_data_scanner_cycle_with_budget(
"Scanner cycle is recovering to a newer durable cache generation" "Scanner cycle is recovering to a newer durable cache generation"
); );
emit_scan_cycle_partial_with_source(cycle_start.elapsed(), ScanCyclePartialReason::Unknown, None); emit_scan_cycle_partial_with_source(cycle_start.elapsed(), ScanCyclePartialReason::Unknown, None);
let persisted = persist_required_scanner_cycle_floor( return if persist_required_scanner_cycle_floor(
ctx, ctx,
storeapi.clone(), storeapi.clone(),
cycle_info, cycle_info,
@@ -1386,9 +1286,8 @@ async fn run_data_scanner_cycle_with_budget(
required_cycle, required_cycle,
&mut cycle_metrics_guard, &mut cycle_metrics_guard,
) )
.await; .await
return if persisted { {
cycle_budget.mark_cycle_state_persisted();
ScannerCycleOutcome::Partial ScannerCycleOutcome::Partial
} else { } else {
ScannerCycleOutcome::Failed ScannerCycleOutcome::Failed
@@ -1446,7 +1345,7 @@ async fn run_data_scanner_cycle_with_budget(
scan_cycle_partial_reason(budget_reason), scan_cycle_partial_reason(budget_reason),
scan_cycle_partial_source(budget_reason), scan_cycle_partial_source(budget_reason),
); );
let persisted = finalize_partial_scan_cycle( return if finalize_partial_scan_cycle(
ctx, ctx,
storeapi.clone(), storeapi.clone(),
cycle_info, cycle_info,
@@ -1454,9 +1353,8 @@ async fn run_data_scanner_cycle_with_budget(
leader_epoch, leader_epoch,
&mut cycle_metrics_guard, &mut cycle_metrics_guard,
) )
.await; .await
return if persisted { {
cycle_budget.mark_cycle_state_persisted();
ScannerCycleOutcome::Partial ScannerCycleOutcome::Partial
} else { } else {
ScannerCycleOutcome::Failed ScannerCycleOutcome::Failed
@@ -1531,7 +1429,7 @@ async fn run_data_scanner_cycle_with_budget(
); );
} }
emit_scan_cycle_partial_with_source(cycle_start.elapsed(), ScanCyclePartialReason::Unknown, None); emit_scan_cycle_partial_with_source(cycle_start.elapsed(), ScanCyclePartialReason::Unknown, None);
let persisted = finalize_partial_scan_cycle( return if finalize_partial_scan_cycle(
ctx, ctx,
storeapi.clone(), storeapi.clone(),
cycle_info, cycle_info,
@@ -1539,9 +1437,8 @@ async fn run_data_scanner_cycle_with_budget(
leader_epoch, leader_epoch,
&mut cycle_metrics_guard, &mut cycle_metrics_guard,
) )
.await; .await
return if persisted { {
cycle_budget.mark_cycle_state_persisted();
ScannerCycleOutcome::Partial ScannerCycleOutcome::Partial
} else { } else {
ScannerCycleOutcome::Failed ScannerCycleOutcome::Failed
@@ -1582,7 +1479,6 @@ async fn run_data_scanner_cycle_with_budget(
) )
.await .await
{ {
cycle_budget.mark_cycle_state_persisted();
emit_scan_cycle_superseded(cycle_start.elapsed()); emit_scan_cycle_superseded(cycle_start.elapsed());
return ScannerCycleOutcome::Superseded; return ScannerCycleOutcome::Superseded;
} }
@@ -1615,7 +1511,6 @@ async fn run_data_scanner_cycle_with_budget(
emit_scan_cycle_complete(false, cycle_start.elapsed()); emit_scan_cycle_complete(false, cycle_start.elapsed());
return ScannerCycleOutcome::Failed; return ScannerCycleOutcome::Failed;
} }
cycle_budget.mark_cycle_state_persisted();
done_cycle(); done_cycle();
emit_scan_cycle_complete(true, cycle_start.elapsed()); emit_scan_cycle_complete(true, cycle_start.elapsed());
@@ -1680,7 +1575,7 @@ async fn run_data_scanner_with_maintenance_state(
) -> Result<(), ScannerError> { ) -> Result<(), ScannerError> {
reset_scanner_cycle_schedule(); reset_scanner_cycle_schedule();
// Acquire leader lock (write lock) to ensure only one scanner runs // Acquire leader lock (write lock) to ensure only one scanner runs
let mut guard = match storeapi.new_ns_lock(RUSTFS_META_BUCKET, "leader.lock").await { let guard = match storeapi.new_ns_lock(RUSTFS_META_BUCKET, "leader.lock").await {
Ok(ns_lock) => match ns_lock.get_write_lock_quiet(get_lock_acquire_timeout()).await { Ok(ns_lock) => match ns_lock.get_write_lock_quiet(get_lock_acquire_timeout()).await {
Ok(guard) => { Ok(guard) => {
record_scanner_leader_lock_state("acquired"); record_scanner_leader_lock_state("acquired");
@@ -1845,49 +1740,13 @@ async fn run_data_scanner_with_maintenance_state(
return Ok(()); return Ok(());
} }
let cycle_ctx = ctx.child_token(); let cycle_ctx = ctx.child_token();
let cycle_budget = ScannerCycleBudget::new_with_runtime_progress_tracking(&cycle_ctx, scanner_cycle_budget_config()); let initial_outcome = await_scanner_cycle_with_lock_fence(
let initial_outcome = match await_scanner_cycle_with_budget_fence(
&cycle_ctx, &cycle_ctx,
&cycle_budget, run_data_scanner_cycle(&cycle_ctx, &storeapi, &mut cycle_info, &mut cycle_revision, leader_epoch),
run_data_scanner_cycle_with_budget(
&cycle_ctx,
&storeapi,
&mut cycle_info,
&mut cycle_revision,
leader_epoch,
cycle_budget.clone(),
),
guard.lock_lost_notified(), guard.lock_lost_notified(),
) )
.await .await
{ .unwrap_or(ScannerCycleOutcome::Failed);
ScannerCycleWaitOutcome::Completed(outcome) => outcome,
ScannerCycleWaitOutcome::LockLost => {
record_scanner_leader_lock_lost("Scanner leader lock lost during the initial cycle").await;
global_metrics().set_cycle(None).await;
return Ok(());
}
ScannerCycleWaitOutcome::Cancelled => {
global_metrics().set_cycle(None).await;
return Ok(());
}
ScannerCycleWaitOutcome::Deadline { worker_stopped } => {
handle_scanner_cycle_deadline(
&ctx,
storeapi.clone(),
ScannerCycleDeadlineState {
cycle_info: &mut cycle_info,
cycle_revision: &mut cycle_revision,
leader_epoch: &mut leader_epoch,
cycle_budget: &cycle_budget,
},
worker_stopped,
&mut guard,
)
.await;
return Ok(());
}
};
superseded_backoff.record_retryable_cycle(initial_outcome == ScannerCycleOutcome::Superseded); superseded_backoff.record_retryable_cycle(initial_outcome == ScannerCycleOutcome::Superseded);
deferred_backoff.record_retryable_cycle(matches!(initial_outcome, ScannerCycleOutcome::Deferred(_))); deferred_backoff.record_retryable_cycle(matches!(initial_outcome, ScannerCycleOutcome::Deferred(_)));
dirty_usage_generation_seen = dirty_generation_before_cycle; dirty_usage_generation_seen = dirty_generation_before_cycle;
@@ -2093,49 +1952,13 @@ async fn run_data_scanner_with_maintenance_state(
} }
let dirty_generation_before_cycle = dirty_usage_generation(); let dirty_generation_before_cycle = dirty_usage_generation();
let cycle_ctx = ctx.child_token(); let cycle_ctx = ctx.child_token();
let cycle_budget = ScannerCycleBudget::new_with_runtime_progress_tracking(&cycle_ctx, scanner_cycle_budget_config()); let outcome = await_scanner_cycle_with_lock_fence(
let outcome = match await_scanner_cycle_with_budget_fence(
&cycle_ctx, &cycle_ctx,
&cycle_budget, run_data_scanner_cycle(&cycle_ctx, &storeapi, &mut cycle_info, &mut cycle_revision, leader_epoch),
run_data_scanner_cycle_with_budget(
&cycle_ctx,
&storeapi,
&mut cycle_info,
&mut cycle_revision,
leader_epoch,
cycle_budget.clone(),
),
guard.lock_lost_notified(), guard.lock_lost_notified(),
) )
.await .await
{ .unwrap_or(ScannerCycleOutcome::Failed);
ScannerCycleWaitOutcome::Completed(outcome) => outcome,
ScannerCycleWaitOutcome::LockLost => {
record_scanner_leader_lock_lost("Scanner leader lock lost during a scanner cycle").await;
global_metrics().set_cycle(None).await;
return Ok(());
}
ScannerCycleWaitOutcome::Cancelled => {
global_metrics().set_cycle(None).await;
return Ok(());
}
ScannerCycleWaitOutcome::Deadline { worker_stopped } => {
handle_scanner_cycle_deadline(
&ctx,
storeapi.clone(),
ScannerCycleDeadlineState {
cycle_info: &mut cycle_info,
cycle_revision: &mut cycle_revision,
leader_epoch: &mut leader_epoch,
cycle_budget: &cycle_budget,
},
worker_stopped,
&mut guard,
)
.await;
return Ok(());
}
};
superseded_backoff.record_retryable_cycle(outcome == ScannerCycleOutcome::Superseded); superseded_backoff.record_retryable_cycle(outcome == ScannerCycleOutcome::Superseded);
deferred_backoff.record_retryable_cycle(matches!(outcome, ScannerCycleOutcome::Deferred(_))); deferred_backoff.record_retryable_cycle(matches!(outcome, ScannerCycleOutcome::Deferred(_)));
dirty_usage_generation_seen = dirty_generation_before_cycle; dirty_usage_generation_seen = dirty_generation_before_cycle;
-60
View File
@@ -1581,63 +1581,3 @@ where
output = &mut cycle => Some(output), output = &mut cycle => Some(output),
} }
} }
#[derive(Debug, PartialEq, Eq)]
pub(super) enum ScannerCycleWaitOutcome<T> {
Completed(T),
LockLost,
Cancelled,
Deadline { worker_stopped: bool },
}
pub(super) async fn await_scanner_cycle_with_budget_fence<Cycle, LockLost>(
cycle_ctx: &CancellationToken,
budget: &ScannerCycleBudget,
cycle: Cycle,
lock_lost: LockLost,
) -> ScannerCycleWaitOutcome<Cycle::Output>
where
Cycle: Future,
LockLost: Future<Output = ()>,
{
tokio::pin!(cycle);
tokio::pin!(lock_lost);
let deadline = async {
if let Some(deadline) = budget.deadline() {
tokio::time::sleep_until(deadline).await;
} else {
std::future::pending::<()>().await;
}
};
tokio::pin!(deadline);
tokio::select! {
biased;
_ = &mut lock_lost => {
cycle_ctx.cancel();
let _ = tokio::time::timeout(SCANNER_LOCK_LOSS_SHUTDOWN_TIMEOUT, &mut cycle).await;
ScannerCycleWaitOutcome::LockLost
}
_ = &mut deadline => {
budget.cancel_for_runtime();
// Let the budget cancellation reach the scanner first so it can
// persist a partial cursor. Only an uncooperative worker gets the
// parent cancellation, and it is dropped after the bounded window;
// the caller fences its epoch next.
let worker_stopped = if tokio::time::timeout(SCANNER_LOCK_LOSS_SHUTDOWN_TIMEOUT, &mut cycle)
.await
.is_ok()
{
true
} else {
cycle_ctx.cancel();
false
};
ScannerCycleWaitOutcome::Deadline { worker_stopped }
}
_ = cycle_ctx.cancelled() => {
let _ = tokio::time::timeout(SCANNER_LOCK_LOSS_SHUTDOWN_TIMEOUT, &mut cycle).await;
ScannerCycleWaitOutcome::Cancelled
}
output = &mut cycle => ScannerCycleWaitOutcome::Completed(output),
}
}
+9 -180
View File
@@ -26,7 +26,6 @@ use std::task::Poll;
use temp_env::{with_var, with_var_unset}; use temp_env::{with_var, with_var_unset};
use tokio::io::AsyncReadExt; use tokio::io::AsyncReadExt;
use tokio::sync::Mutex; use tokio::sync::Mutex;
use tokio::time::{Duration, advance};
const TEST_DEFAULT_SCANNER_CYCLE_SECS: u64 = 24 * 60 * 60; const TEST_DEFAULT_SCANNER_CYCLE_SECS: u64 = 24 * 60 * 60;
@@ -119,178 +118,6 @@ async fn scanner_cycle_lock_fence_bounds_uncooperative_shutdown() {
assert!(cycle_ctx.is_cancelled()); assert!(cycle_ctx.is_cancelled());
} }
#[tokio::test(start_paused = true)]
async fn cycle_budget_fences_late_writer_after_timeout() {
let cycle_ctx = CancellationToken::new();
let budget = ScannerCycleBudget::new(
&cycle_ctx,
ScannerCycleBudgetConfig {
max_duration: Some(Duration::from_secs(5)),
..Default::default()
},
);
let outcome = {
let cycle = std::future::pending::<()>();
let lock_lost = std::future::pending::<()>();
let waiter = await_scanner_cycle_with_budget_fence(&cycle_ctx, &budget, cycle, lock_lost);
tokio::pin!(waiter);
tokio::task::yield_now().await;
advance(Duration::from_secs(5)).await;
tokio::task::yield_now().await;
advance(SCANNER_LOCK_LOSS_SHUTDOWN_TIMEOUT).await;
waiter.await
};
assert_eq!(outcome, ScannerCycleWaitOutcome::Deadline { worker_stopped: false });
assert!(cycle_ctx.is_cancelled());
assert_eq!(budget.reason(), Some(ScannerCycleBudgetReason::Runtime));
// A newer leadership epoch is the durable fence that rejects a late
// writer after the timed-out future has been dropped.
let store = Arc::new(MemoryConfigStore::default());
let mut revision = DataUsageCacheRevision::Missing;
let mut cycle = CurrentCycle {
current: 0,
next: 12,
..Default::default()
};
let persist_ctx = CancellationToken::new();
assert!(persist_scanner_cycle_state(&persist_ctx, store.clone(), &mut cycle, &mut revision, 1).await);
let newer = encode_scanner_cycle_state(&cycle, 2).expect("new epoch fence should encode");
let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
store.interleaving_puts.lock().await.insert(key, (2, newer));
let mut late_cycle = CurrentCycle { next: 13, ..cycle };
assert!(!persist_scanner_cycle_state(&persist_ctx, store, &mut late_cycle, &mut revision, 1).await);
}
#[tokio::test(start_paused = true)]
async fn cycle_budget_parent_cancellation_is_not_reported_as_timeout() {
let cycle_ctx = CancellationToken::new();
let budget = ScannerCycleBudget::new(
&cycle_ctx,
ScannerCycleBudgetConfig {
max_duration: Some(Duration::from_secs(5)),
..Default::default()
},
);
let waiter = await_scanner_cycle_with_budget_fence(&cycle_ctx, &budget, std::future::pending::<()>(), std::future::pending());
tokio::pin!(waiter);
tokio::task::yield_now().await;
cycle_ctx.cancel();
tokio::task::yield_now().await;
advance(SCANNER_LOCK_LOSS_SHUTDOWN_TIMEOUT).await;
assert_eq!(waiter.await, ScannerCycleWaitOutcome::Cancelled);
}
#[tokio::test(start_paused = true)]
async fn cycle_budget_deadline_wins_same_tick_as_parent_cancellation() {
let cycle_ctx = CancellationToken::new();
let budget = ScannerCycleBudget::new(
&cycle_ctx,
ScannerCycleBudgetConfig {
max_duration: Some(Duration::from_secs(5)),
..Default::default()
},
);
let waiter = await_scanner_cycle_with_budget_fence(&cycle_ctx, &budget, std::future::pending::<()>(), std::future::pending());
tokio::pin!(waiter);
tokio::task::yield_now().await;
advance(Duration::from_secs(5)).await;
cycle_ctx.cancel();
tokio::task::yield_now().await;
advance(SCANNER_LOCK_LOSS_SHUTDOWN_TIMEOUT).await;
assert_eq!(waiter.await, ScannerCycleWaitOutcome::Deadline { worker_stopped: false });
assert_eq!(budget.reason(), Some(ScannerCycleBudgetReason::Runtime));
}
#[tokio::test]
async fn cycle_budget_persist_cursor_failure_is_recovery_required() {
let store = Arc::new(MemoryConfigStore::default());
let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
store.fail_put_number.lock().await.insert(key, 1);
let ctx = CancellationToken::new();
let mut revision = DataUsageCacheRevision::Missing;
let mut cycle = CurrentCycle {
current: 12,
next: 12,
..Default::default()
};
let mut leader_epoch = 1;
let fenced = fence_scanner_epoch_after_cycle_timeout(
&ctx,
store,
&mut cycle,
&mut revision,
&mut leader_epoch,
std::future::pending(),
)
.await;
assert!(!fenced, "a failed cursor/generation write must require recovery");
let budget = ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default());
assert!(cycle_timeout_requires_recovery(true, budget.cycle_state_persisted(), fenced));
let metrics = Metrics::new();
metrics.record_scanner_cycle_timeout(!fenced, Duration::from_secs(17));
let report = metrics.report().await;
assert_eq!(report.cycle_timeout_total, 1);
assert_eq!(report.cycle_recovery_required_total, 1);
assert_eq!(report.cycle_last_progress_age, 17);
assert!(report.leader_lease_without_progress);
}
#[tokio::test]
async fn cycle_budget_deadline_handler_fences_and_releases_guard() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
let lock = store
.new_ns_lock(RUSTFS_META_BUCKET, "leader.lock")
.await
.expect("scanner leader lock should be created");
let mut guard = lock
.get_write_lock(Duration::from_secs(1))
.await
.expect("scanner leader lock should be acquired");
let ctx = CancellationToken::new();
let mut cycle_info = CurrentCycle {
current: 12,
next: 12,
..Default::default()
};
let mut cycle_revision = DataUsageCacheRevision::Missing;
let mut leader_epoch = 1;
let budget = ScannerCycleBudget::new(
&ctx,
ScannerCycleBudgetConfig {
max_duration: Some(Duration::from_secs(60)),
..Default::default()
},
);
budget.mark_cycle_state_persisted();
handle_scanner_cycle_deadline(
&ctx,
store.clone(),
ScannerCycleDeadlineState {
cycle_info: &mut cycle_info,
cycle_revision: &mut cycle_revision,
leader_epoch: &mut leader_epoch,
cycle_budget: &budget,
},
true,
&mut guard,
)
.await;
assert!(guard.is_released());
let persisted = read_config(store, &DATA_USAGE_BLOOM_NAME_PATH)
.await
.expect("deadline handler should persist a fenced cursor");
let (_, persisted_epoch) = decode_scanner_cycle_state(&persisted).expect("fenced cursor should decode");
assert_eq!(persisted_epoch, 2);
global_metrics().set_cycle(None).await;
}
#[tokio::test] #[tokio::test]
async fn scanner_cycle_recovery_wake_survives_wait_registration_race() { async fn scanner_cycle_recovery_wake_survives_wait_registration_race() {
notify_scanner_cycle_recovery_wake(); notify_scanner_cycle_recovery_wake();
@@ -601,6 +428,13 @@ fn test_scanner_cycle_max_duration_uses_env() {
}); });
} }
#[test]
fn test_scanner_cycle_max_duration_default_is_disabled() {
with_var_unset(ENV_SCANNER_CYCLE_MAX_DURATION_SECS, || {
assert_eq!(scanner_cycle_max_duration(), None);
});
}
#[tokio::test] #[tokio::test]
async fn test_scanner_cycle_budget_cancels_after_duration() { async fn test_scanner_cycle_budget_cancels_after_duration() {
let parent = CancellationToken::new(); let parent = CancellationToken::new();
@@ -2408,7 +2242,7 @@ async fn test_leadership_claim_usage_fence_rejects_old_inflight_writer() {
} }
#[tokio::test] #[tokio::test]
async fn cycle_budget_lease_takeover_rejects_old_generation() { async fn test_successful_old_epoch_commit_is_fenced_after_cancellation() {
let store = Arc::new(MemoryConfigStore::default()); let store = Arc::new(MemoryConfigStore::default());
let ctx = CancellationToken::new(); let ctx = CancellationToken::new();
let mut revision = DataUsageCacheRevision::Missing; let mut revision = DataUsageCacheRevision::Missing;
@@ -2453,17 +2287,12 @@ async fn cycle_budget_lease_takeover_rejects_old_generation() {
.await .await
); );
let state = read_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH) let state = read_config(store, &DATA_USAGE_BLOOM_NAME_PATH)
.await .await
.expect("replacement leadership claim should persist"); .expect("replacement leadership claim should persist");
let (claimed_cycle, claimed_epoch) = decode_scanner_cycle_state(&state).expect("replacement cycle state should decode"); let (claimed_cycle, claimed_epoch) = decode_scanner_cycle_state(&state).expect("replacement cycle state should decode");
assert_eq!(claimed_cycle.next, 14); assert_eq!(claimed_cycle.next, 14);
assert_eq!(claimed_epoch, 2); assert_eq!(claimed_epoch, 2);
let mut stale_cycle = CurrentCycle { next: 15, ..cycle };
let mut stale_revision = DataUsageCacheRevision::Etag("memory-2".to_string());
let stale_ctx = CancellationToken::new();
assert!(!persist_scanner_cycle_state(&stale_ctx, store, &mut stale_cycle, &mut stale_revision, 1,).await);
} }
#[tokio::test] #[tokio::test]
+15 -146
View File
@@ -14,16 +14,17 @@
use std::sync::{ use std::sync::{
Arc, Arc,
atomic::{AtomicBool, AtomicU8, AtomicU64, Ordering}, atomic::{AtomicU8, AtomicU64, Ordering},
}; };
use tokio::time::{Duration, Instant}; use std::time::Instant;
use tokio::time::Duration;
use tokio_util::sync::CancellationToken; use tokio_util::sync::CancellationToken;
const BUDGET_REASON_NONE: u8 = 0; const BUDGET_REASON_NONE: u8 = 0;
const BUDGET_REASON_RUNTIME: u8 = 1; const BUDGET_REASON_RUNTIME: u8 = 1;
const BUDGET_REASON_OBJECTS: u8 = 2; const BUDGET_REASON_OBJECTS: u8 = 2;
const BUDGET_REASON_DIRECTORIES: u8 = 3; const BUDGET_REASON_DIRECTORIES: u8 = 3;
const PROGRESS_CLOCK_SAMPLE_INTERVAL: u64 = 128;
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] #[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub(crate) struct ScannerCycleBudgetConfig { pub(crate) struct ScannerCycleBudgetConfig {
@@ -62,51 +63,29 @@ pub struct ScannerCycleBudget {
token: CancellationToken, token: CancellationToken,
reason: Arc<AtomicU8>, reason: Arc<AtomicU8>,
started_at: Instant, started_at: Instant,
deadline: Option<Instant>,
max_duration: Option<Duration>, max_duration: Option<Duration>,
max_objects: Option<u64>, max_objects: Option<u64>,
max_directories: Option<u64>, max_directories: Option<u64>,
track_progress: bool, track_progress: bool,
track_unbounded_counts: bool,
objects_scanned: AtomicU64, objects_scanned: AtomicU64,
directories_started: AtomicU64, directories_started: AtomicU64,
entries_visited: AtomicU64, entries_visited: AtomicU64,
last_progress_millis: AtomicU64,
cycle_state_persisted: AtomicBool,
} }
impl ScannerCycleBudget { impl ScannerCycleBudget {
#[cfg(test)]
pub(crate) fn new(parent: &CancellationToken, config: ScannerCycleBudgetConfig) -> Arc<Self> { pub(crate) fn new(parent: &CancellationToken, config: ScannerCycleBudgetConfig) -> Arc<Self> {
Self::new_inner(parent, config, false, false) Self::new_inner(parent, config, false)
} }
pub(crate) fn new_with_progress_tracking(parent: &CancellationToken, config: ScannerCycleBudgetConfig) -> Arc<Self> { pub(crate) fn new_with_progress_tracking(parent: &CancellationToken, config: ScannerCycleBudgetConfig) -> Arc<Self> {
Self::new_inner(parent, config, true, true) Self::new_inner(parent, config, true)
} }
pub(crate) fn new_with_runtime_progress_tracking(parent: &CancellationToken, config: ScannerCycleBudgetConfig) -> Arc<Self> { fn new_inner(parent: &CancellationToken, config: ScannerCycleBudgetConfig, track_progress: bool) -> Arc<Self> {
let track_progress = config.max_duration.is_some();
Self::new_inner(parent, config, track_progress, false)
}
fn new_inner(
parent: &CancellationToken,
config: ScannerCycleBudgetConfig,
track_progress: bool,
track_unbounded_counts: bool,
) -> Arc<Self> {
let token = parent.child_token(); let token = parent.child_token();
let reason = Arc::new(AtomicU8::new(BUDGET_REASON_NONE)); let reason = Arc::new(AtomicU8::new(BUDGET_REASON_NONE));
let started_at = Instant::now();
let deadline = config.max_duration.map(|duration| match started_at.checked_add(duration) {
Some(deadline) => deadline,
// Runtime config rejects this range, but keep programmatic callers
// fail-closed instead of panicking or silently disabling the wall clock.
None => started_at,
});
if let Some(deadline) = deadline { if let Some(duration) = config.max_duration {
let parent = parent.clone(); let parent = parent.clone();
let token_wait = token.clone(); let token_wait = token.clone();
let token_cancel = token.clone(); let token_cancel = token.clone();
@@ -115,7 +94,7 @@ impl ScannerCycleBudget {
tokio::select! { tokio::select! {
_ = parent.cancelled() => {} _ = parent.cancelled() => {}
_ = token_wait.cancelled() => {} _ = token_wait.cancelled() => {}
_ = tokio::time::sleep_until(deadline) => { _ = tokio::time::sleep(duration) => {
Self::cancel_for_reason(&reason, &token_cancel, ScannerCycleBudgetReason::Runtime); Self::cancel_for_reason(&reason, &token_cancel, ScannerCycleBudgetReason::Runtime);
} }
} }
@@ -125,18 +104,14 @@ impl ScannerCycleBudget {
Arc::new(Self { Arc::new(Self {
token, token,
reason, reason,
started_at, started_at: Instant::now(),
deadline,
max_duration: config.max_duration, max_duration: config.max_duration,
max_objects: config.max_objects, max_objects: config.max_objects,
max_directories: config.max_directories, max_directories: config.max_directories,
track_progress, track_progress,
track_unbounded_counts,
objects_scanned: AtomicU64::new(0), objects_scanned: AtomicU64::new(0),
directories_started: AtomicU64::new(0), directories_started: AtomicU64::new(0),
entries_visited: AtomicU64::new(0), entries_visited: AtomicU64::new(0),
last_progress_millis: AtomicU64::new(0),
cycle_state_persisted: AtomicBool::new(false),
}) })
} }
@@ -156,14 +131,6 @@ impl ScannerCycleBudget {
self.max_duration self.max_duration
} }
pub(crate) fn deadline(&self) -> Option<Instant> {
self.deadline
}
pub(crate) fn cancel_for_runtime(&self) {
self.cancel_for(ScannerCycleBudgetReason::Runtime);
}
pub(crate) fn max_objects(&self) -> Option<u64> { pub(crate) fn max_objects(&self) -> Option<u64> {
self.max_objects self.max_objects
} }
@@ -206,43 +173,15 @@ impl ScannerCycleBudget {
self.entries_visited.load(Ordering::Relaxed) self.entries_visited.load(Ordering::Relaxed)
} }
pub(crate) fn mark_cycle_state_persisted(&self) {
self.cycle_state_persisted.store(true, Ordering::Release);
}
pub(crate) fn cycle_state_persisted(&self) -> bool {
self.cycle_state_persisted.load(Ordering::Acquire)
}
pub(crate) fn progress_age(&self) -> Duration {
let elapsed_millis = u64::try_from(self.started_at.elapsed().as_millis()).unwrap_or(u64::MAX);
let last_progress = self.last_progress_millis.load(Ordering::Relaxed);
Duration::from_millis(elapsed_millis.saturating_sub(last_progress))
}
fn record_progress_sample(&self, event: u64) {
// Clock reads are sampled at batch/count boundaries; the scanner's
// per-object path does not add a second progress atomic.
if event == 0 || (event != 1 && !event.is_multiple_of(PROGRESS_CLOCK_SAMPLE_INTERVAL)) {
return;
}
let elapsed_millis = u64::try_from(self.started_at.elapsed().as_millis()).unwrap_or(u64::MAX);
self.last_progress_millis.store(elapsed_millis, Ordering::Relaxed);
}
pub(crate) fn record_entries_visited(&self, entries_visited: u64) { pub(crate) fn record_entries_visited(&self, entries_visited: u64) {
if self.track_progress { if self.track_progress {
let entries = saturating_fetch_add(&self.entries_visited, entries_visited); saturating_fetch_add(&self.entries_visited, entries_visited);
self.record_progress_sample(entries);
} }
} }
pub(crate) fn record_remote_progress(&self, objects_scanned: u64, directories_started: u64) { pub(crate) fn record_remote_progress(&self, objects_scanned: u64, directories_started: u64) {
if self.track_progress || self.max_objects.is_some() { if self.track_progress || self.max_objects.is_some() {
let objects = saturating_fetch_add(&self.objects_scanned, objects_scanned); let objects = saturating_fetch_add(&self.objects_scanned, objects_scanned);
if self.track_progress {
self.record_progress_sample(objects);
}
if self.max_objects.is_some_and(|max_objects| objects >= max_objects) { if self.max_objects.is_some_and(|max_objects| objects >= max_objects) {
self.cancel_for(ScannerCycleBudgetReason::Objects); self.cancel_for(ScannerCycleBudgetReason::Objects);
} }
@@ -250,12 +189,9 @@ impl ScannerCycleBudget {
if self.track_progress || self.max_directories.is_some() { if self.track_progress || self.max_directories.is_some() {
let directories = saturating_fetch_add(&self.directories_started, directories_started); let directories = saturating_fetch_add(&self.directories_started, directories_started);
if self.track_progress {
self.record_progress_sample(directories);
}
if self if self
.max_directories .max_directories
.is_some_and(|max_directories| directory_budget_exhausted(directories, max_directories)) .is_some_and(|max_directories| directories > max_directories)
{ {
self.cancel_for(ScannerCycleBudgetReason::Directories); self.cancel_for(ScannerCycleBudgetReason::Directories);
} }
@@ -271,17 +207,14 @@ impl ScannerCycleBudget {
} }
pub(crate) fn try_start_directory(&self) -> bool { pub(crate) fn try_start_directory(&self) -> bool {
if self.max_directories.is_none() && !self.track_unbounded_counts { if !self.track_progress && self.max_directories.is_none() {
return true; return true;
} }
let directories = saturating_fetch_add(&self.directories_started, 1); let directories = saturating_fetch_add(&self.directories_started, 1);
if self.track_progress {
self.record_progress_sample(directories);
}
if self if self
.max_directories .max_directories
.is_some_and(|max_directories| directory_budget_exhausted(directories, max_directories)) .is_some_and(|max_directories| directories > max_directories)
{ {
self.cancel_for(ScannerCycleBudgetReason::Directories); self.cancel_for(ScannerCycleBudgetReason::Directories);
return false; return false;
@@ -291,14 +224,11 @@ impl ScannerCycleBudget {
} }
pub(crate) fn record_object_scanned(&self) { pub(crate) fn record_object_scanned(&self) {
if self.max_objects.is_none() && !self.track_unbounded_counts { if !self.track_progress && self.max_objects.is_none() {
return; return;
} }
let objects = saturating_fetch_add(&self.objects_scanned, 1); let objects = saturating_fetch_add(&self.objects_scanned, 1);
if self.track_progress {
self.record_progress_sample(objects);
}
if self.max_objects.is_some_and(|max_objects| objects >= max_objects) { if self.max_objects.is_some_and(|max_objects| objects >= max_objects) {
self.cancel_for(ScannerCycleBudgetReason::Objects); self.cancel_for(ScannerCycleBudgetReason::Objects);
} }
@@ -329,13 +259,6 @@ fn saturating_fetch_add(value: &AtomicU64, delta: u64) -> u64 {
} }
} }
fn directory_budget_exhausted(directories: u64, max_directories: u64) -> bool {
// Saturation hides a remote max+1 update when the configured limit is the
// largest representable counter. Treat that boundary as exhausted rather
// than allowing work to continue indefinitely.
directories > max_directories || (directories == u64::MAX && max_directories == u64::MAX)
}
impl Drop for ScannerCycleBudget { impl Drop for ScannerCycleBudget {
fn drop(&mut self) { fn drop(&mut self) {
self.token.cancel(); self.token.cancel();
@@ -478,35 +401,6 @@ mod tests {
assert_eq!(directory_budget.reason(), Some(ScannerCycleBudgetReason::Directories)); assert_eq!(directory_budget.reason(), Some(ScannerCycleBudgetReason::Directories));
} }
#[test]
fn directory_budget_fails_closed_when_progress_saturates() {
let parent = CancellationToken::new();
let budget = ScannerCycleBudget::new(
&parent,
ScannerCycleBudgetConfig {
max_directories: Some(u64::MAX),
..Default::default()
},
);
budget.record_remote_progress(0, u64::MAX);
assert_eq!(budget.reason(), Some(ScannerCycleBudgetReason::Directories));
assert!(budget.token().is_cancelled());
let local_budget = ScannerCycleBudget::new(
&parent,
ScannerCycleBudgetConfig {
max_directories: Some(u64::MAX),
..Default::default()
},
);
local_budget.record_remote_progress(0, u64::MAX - 1);
assert!(!local_budget.budget_elapsed());
assert!(!local_budget.try_start_directory());
assert_eq!(local_budget.reason(), Some(ScannerCycleBudgetReason::Directories));
}
#[test] #[test]
fn explicit_progress_tracking_counts_unbounded_remote_work_without_cancelling() { fn explicit_progress_tracking_counts_unbounded_remote_work_without_cancelling() {
let parent = CancellationToken::new(); let parent = CancellationToken::new();
@@ -567,29 +461,4 @@ mod tests {
assert!(object_limited.requires_serial_progress_accounting()); assert!(object_limited.requires_serial_progress_accounting());
assert!(directory_limited.requires_serial_progress_accounting()); assert!(directory_limited.requires_serial_progress_accounting());
} }
#[tokio::test(start_paused = true)]
async fn progress_age_uses_virtual_time_and_sampled_progress() {
let parent = CancellationToken::new();
let budget = ScannerCycleBudget::new_with_runtime_progress_tracking(
&parent,
ScannerCycleBudgetConfig {
max_duration: Some(Duration::from_secs(60)),
..Default::default()
},
);
tokio::time::advance(Duration::from_secs(5)).await;
assert_eq!(budget.progress_age(), Duration::from_secs(5));
budget.record_entries_visited(1);
assert_eq!(budget.progress_age(), Duration::ZERO);
tokio::time::advance(Duration::from_secs(2)).await;
for _ in 0..126 {
budget.record_entries_visited(1);
}
assert_eq!(budget.progress_age(), Duration::from_secs(2));
budget.record_entries_visited(1);
assert_eq!(budget.progress_age(), Duration::ZERO);
}
} }
-1
View File
@@ -48,7 +48,6 @@ use time::OffsetDateTime;
use tokio::sync::{Mutex, Notify, Semaphore, mpsc}; use tokio::sync::{Mutex, Notify, Semaphore, mpsc};
use tokio::time::Duration; use tokio::time::Duration;
use tokio_util::sync::CancellationToken; use tokio_util::sync::CancellationToken;
use tokio_util::task::AbortOnDropHandle;
use tracing::{debug, error, warn}; use tracing::{debug, error, warn};
use crate::ScannerObjectInfo as ObjectInfo; use crate::ScannerObjectInfo as ObjectInfo;
+4 -4
View File
@@ -314,7 +314,7 @@ impl ScannerIOCache for SetDisks {
let ctx_clone = ctx.clone(); let ctx_clone = ctx.clone();
let completed_bucket_count = Arc::new(AtomicUsize::new(0)); let completed_bucket_count = Arc::new(AtomicUsize::new(0));
let completed_bucket_count_clone = completed_bucket_count.clone(); let completed_bucket_count_clone = completed_bucket_count.clone();
let collect_bucket_results_fut = AbortOnDropHandle::new(tokio::spawn(async move { let collect_bucket_results_fut = tokio::spawn(async move {
let mut cancelled = false; let mut cancelled = false;
loop { loop {
@@ -333,7 +333,7 @@ impl ScannerIOCache for SetDisks {
} }
} }
} }
})); });
let mut futs = Vec::new(); let mut futs = Vec::new();
@@ -365,7 +365,7 @@ impl ScannerIOCache for SetDisks {
NamespaceScannerWorkerMode::RemoteV4(server_epoch) => Some(server_epoch), NamespaceScannerWorkerMode::RemoteV4(server_epoch) => Some(server_epoch),
NamespaceScannerWorkerMode::Coordinator => None, NamespaceScannerWorkerMode::Coordinator => None,
}; };
futs.push(AbortOnDropHandle::new(tokio::spawn(async move { futs.push(tokio::spawn(async move {
let remote_session_id = uuid::Uuid::new_v4(); let remote_session_id = uuid::Uuid::new_v4();
let mut remote_session_sequence = 0_u64; let mut remote_session_sequence = 0_u64;
loop { loop {
@@ -1038,7 +1038,7 @@ impl ScannerIOCache for SetDisks {
); );
} }
} }
}))); }));
} }
drop(bucket_tx); drop(bucket_tx);
drop(bucket_result_tx); drop(bucket_result_tx);
+2 -2
View File
@@ -242,7 +242,7 @@ impl ScannerIOCycle for ECStore {
results[results_index_clone] = result; results[results_index_clone] = result;
} }
}); });
wait_futs.push(AbortOnDropHandle::new(receiver_fut)); wait_futs.push(receiver_fut);
let scan_plan = ScannerBucketScanPlan { let scan_plan = ScannerBucketScanPlan {
buckets: set_buckets, buckets: set_buckets,
@@ -318,7 +318,7 @@ impl ScannerIOCycle for ECStore {
record_set_scan_failure(&mut first_err, e); record_set_scan_failure(&mut first_err, e);
} }
}); });
wait_futs.push(AbortOnDropHandle::new(scanner_fut)); wait_futs.push(scanner_fut);
} }
} }
+1
View File
@@ -76,6 +76,7 @@ pub use bucket::{BucketInfo, BucketOperations, BucketOptions, DeleteBucketOption
pub use capability::{CapabilitySnapshotError, CapabilityState, CapabilityStatus}; pub use capability::{CapabilitySnapshotError, CapabilityState, CapabilityStatus};
pub use error::{StorageErrorCode, StorageResult}; pub use error::{StorageErrorCode, StorageResult};
pub use multipart::{CompletePart, ListMultipartsInfo, ListPartsInfo, MultipartInfo, MultipartUploadResult, PartInfo}; pub use multipart::{CompletePart, ListMultipartsInfo, ListPartsInfo, MultipartInfo, MultipartUploadResult, PartInfo};
pub use object::DeleteAccounting;
pub use object::ObjectLockDeleteOptions; pub use object::ObjectLockDeleteOptions;
pub use object::{DeletedObject, ObjectToDelete}; pub use object::{DeletedObject, ObjectToDelete};
pub use object::{ExpirationOptions, TransitionedObject}; pub use object::{ExpirationOptions, TransitionedObject};
+24
View File
@@ -218,6 +218,17 @@ pub struct DeletedObject {
pub force_delete_generation: Option<i64>, pub force_delete_generation: Option<i64>,
} }
/// Accounting identity returned by the internal commit-time delete path.
///
/// This is carried separately from [`DeletedObject`] so adding quota details
/// does not change the source shape of the public S3 delete result contract.
#[derive(Debug, Default, Clone, PartialEq, Eq)]
pub struct DeleteAccounting {
pub size: Option<u64>,
pub version_id: Option<Uuid>,
pub removed_current_object: bool,
}
impl DeletedObject { impl DeletedObject {
pub fn version_purge_status(&self) -> VersionPurgeStatusType { pub fn version_purge_status(&self) -> VersionPurgeStatusType {
self.replication_state self.replication_state
@@ -341,6 +352,19 @@ pub trait ObjectOperations: Send + Sync + fmt::Debug {
objects: Vec<Self::ObjectToDelete>, objects: Vec<Self::ObjectToDelete>,
opts: Self::ObjectOptions, opts: Self::ObjectOptions,
) -> (Vec<Self::DeletedObject>, Vec<Option<Self::Error>>); ) -> (Vec<Self::DeletedObject>, Vec<Option<Self::Error>>);
/// Delete objects and optionally return commit-time accounting identities.
/// The default preserves the ordinary delete contract for implementations
/// that do not expose storage-level accounting details.
async fn delete_objects_with_accounting(
&self,
bucket: &str,
objects: Vec<Self::ObjectToDelete>,
opts: Self::ObjectOptions,
) -> (Vec<Self::DeletedObject>, Vec<Option<Self::Error>>, Vec<Option<DeleteAccounting>>) {
let object_count = objects.len();
let (deleted, errors) = self.delete_objects(bucket, objects, opts).await;
(deleted, errors, vec![None; object_count])
}
async fn put_object_metadata( async fn put_object_metadata(
&self, &self,
bucket: &str, bucket: &str,
+2 -2
View File
@@ -268,7 +268,7 @@ where
.parse::<T>() .parse::<T>()
.map_err(|_| { .map_err(|_| {
log_once(&format!("env_invalid_value:{used_key}"), || { log_once(&format!("env_invalid_value:{used_key}"), || {
format!("Invalid {} value for {used_key}. Treating as unset.", type_name::<T>()) format!("Invalid {} value for {used_key}: {value}. Treating as unset.", type_name::<T>())
}); });
}) })
.ok() .ok()
@@ -570,7 +570,7 @@ where
Ok(parsed) => EnvParseOutcome::Parsed(parsed), Ok(parsed) => EnvParseOutcome::Parsed(parsed),
Err(_) => { Err(_) => {
log_once(&format!("env_invalid_value:{used_key}"), || { log_once(&format!("env_invalid_value:{used_key}"), || {
format!("Invalid {} value for {used_key}. Treating as unset.", type_name::<T>()) format!("Invalid {} value for {used_key}: {value}. Treating as unset.", type_name::<T>())
}); });
EnvParseOutcome::Invalid EnvParseOutcome::Invalid
} }
+1 -20
View File
@@ -52,7 +52,7 @@ The `/v3/scanner/status` response reports each effective runtime value with a
| `scanner.max_wait` | `RUSTFS_SCANNER_MAX_WAIT_SECS` | seconds | preset-derived | Caps one scanner sleep. | | `scanner.max_wait` | `RUSTFS_SCANNER_MAX_WAIT_SECS` | seconds | preset-derived | Caps one scanner sleep. |
| `scanner.cycle` | `RUSTFS_SCANNER_CYCLE` | seconds | preset-derived | Sets the interval between scanner cycles. | | `scanner.cycle` | `RUSTFS_SCANNER_CYCLE` | seconds | preset-derived | Sets the interval between scanner cycles. |
| `scanner.start_delay` | `RUSTFS_SCANNER_START_DELAY_SECS` | seconds | unset | Sets startup delay and, for compatibility, the cycle interval when `scanner.cycle` is unset. | | `scanner.start_delay` | `RUSTFS_SCANNER_START_DELAY_SECS` | seconds | unset | Sets startup delay and, for compatibility, the cycle interval when `scanner.cycle` is unset. |
| `scanner.cycle_max_duration` | `RUSTFS_SCANNER_CYCLE_MAX_DURATION_SECS` | seconds | `1800` | Caps one cycle's runtime. An explicit `0` disables this budget. | | `scanner.cycle_max_duration` | `RUSTFS_SCANNER_CYCLE_MAX_DURATION_SECS` | seconds | `0` | Caps one cycle's runtime. `0` disables this budget. |
| `scanner.cycle_max_objects` | `RUSTFS_SCANNER_CYCLE_MAX_OBJECTS` | objects | `0` | Caps objects processed by one cycle. `0` disables this budget. | | `scanner.cycle_max_objects` | `RUSTFS_SCANNER_CYCLE_MAX_OBJECTS` | objects | `0` | Caps objects processed by one cycle. `0` disables this budget. |
| `scanner.cycle_max_directories` | `RUSTFS_SCANNER_CYCLE_MAX_DIRECTORIES` | directories | `0` | Caps directories entered by one cycle. `0` disables this budget. | | `scanner.cycle_max_directories` | `RUSTFS_SCANNER_CYCLE_MAX_DIRECTORIES` | directories | `0` | Caps directories entered by one cycle. `0` disables this budget. |
| `heal.bitrot_cycle` | `RUSTFS_SCANNER_BITROT_CYCLE_SECS` | seconds | `2592000` | Controls periodic deep bitrot scans. `false`, `off`, `no`, or `disabled` disables periodic deep scans; `0`, `true`, `on`, or `yes` runs deep mode every scanner cycle. | | `heal.bitrot_cycle` | `RUSTFS_SCANNER_BITROT_CYCLE_SECS` | seconds | `2592000` | Controls periodic deep bitrot scans. `false`, `off`, `no`, or `disabled` disables periodic deep scans; `0`, `true`, `on`, or `yes` runs deep mode every scanner cycle. |
@@ -70,21 +70,6 @@ sleep multiplier, maximum wait, and cycle interval. Use `scanner.delay`,
`scanner.max_wait`, and `scanner.cycle` when the preset is close but one axis `scanner.max_wait`, and `scanner.cycle` when the preset is close but one axis
needs a precise override. needs a precise override.
When the cycle duration control is unset, RustFS uses a finite 1800-second
(30-minute) default, matching the scanner benchmark guidance. An explicit `0`
preserves the compatibility behavior of an unbounded cycle; object and
directory budgets likewise remain unbounded when explicitly set to `0`. Invalid
or overflowing duration environment values are configuration errors rather than
silent fallback values.
When a finite deadline expires, RustFS cancels cooperative scanner work and
waits only for the existing bounded shutdown window. A non-yielding I/O future
is dropped after that window. RustFS then attempts a higher leadership epoch so
late cycle, usage, cache, and remote writes from the old generation fail closed.
If the worker cannot stop cooperatively, the cycle state was not confirmed
durable, or that epoch fence cannot be durably persisted, the scanner reports
`recovery-required`; it does not claim an uncooperative cursor was saved.
An explicit `scanner.cycle` or `RUSTFS_SCANNER_CYCLE` is a minimum inter-cycle An explicit `scanner.cycle` or `RUSTFS_SCANNER_CYCLE` is a minimum inter-cycle
cadence: dirty-usage notifications do not bypass that configured interval. cadence: dirty-usage notifications do not bypass that configured interval.
The default adaptive policy continues to use dirty-usage notifications to wake The default adaptive policy continues to use dirty-usage notifications to wake
@@ -159,10 +144,6 @@ metrics.maintenance_control.primary_control
metrics.source_work metrics.source_work
metrics.replication_repair metrics.replication_repair
metrics.scan_checkpoint metrics.scan_checkpoint
metrics.cycle_timeout_total
metrics.cycle_last_progress_age
metrics.leader_lease_without_progress
metrics.cycle_recovery_required_total
``` ```
## Reading Pacing Pressure ## Reading Pacing Pressure
+1 -1
View File
@@ -1021,7 +1021,7 @@ fn build_list_objects_v2_metadata_output(
object: Object { object: Object {
key: Some(encode_list_objects_v2_value(&object.name, encoding_type)), key: Some(encode_list_objects_v2_value(&object.name, encoding_type)),
last_modified: object.mod_time.map(Timestamp::from), last_modified: object.mod_time.map(Timestamp::from),
size: Some(object.get_actual_size().unwrap_or_default()), size: Some(object.get_actual_size_or_physical()),
e_tag: object.etag.clone().map(|etag| to_s3s_etag(&etag)), e_tag: object.etag.clone().map(|etag| to_s3s_etag(&etag)),
storage_class: Some(ObjectStorageClass::from( storage_class: Some(ObjectStorageClass::from(
object object
+234 -37
View File
@@ -3969,6 +3969,55 @@ fn delete_creates_delete_marker(opts: &ObjectOptions) -> bool {
opts.version_id.is_none() && opts.versioned && !opts.version_suspended opts.version_id.is_none() && opts.versioned && !opts.version_suspended
} }
fn delete_removes_current_object(opts: &ObjectOptions) -> bool {
delete_request_targets_current(
opts.version_id
.as_deref()
.and_then(|version_id| Uuid::parse_str(version_id).ok()),
)
}
fn delete_request_targets_current(version_id: Option<Uuid>) -> bool {
version_id.is_none() || version_id.is_some_and(|version_id| version_id.is_nil())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum DeleteMemoryUpdate {
DeleteMarker,
Object { size: u64, removed_current_object: bool },
}
fn delete_memory_update(
creates_delete_marker: bool,
committed_delete_marker: bool,
requested_current: bool,
accounting_size: Option<u64>,
removed_current_object: bool,
) -> Option<DeleteMemoryUpdate> {
if creates_delete_marker || (committed_delete_marker && requested_current) {
return Some(DeleteMemoryUpdate::DeleteMarker);
}
(!committed_delete_marker)
.then_some(accounting_size)
.flatten()
.map(|size| DeleteMemoryUpdate::Object {
size,
removed_current_object,
})
}
async fn apply_delete_memory_update(bucket: &str, update: Option<DeleteMemoryUpdate>) {
match update {
Some(DeleteMemoryUpdate::DeleteMarker) => record_bucket_delete_marker_memory(bucket).await,
Some(DeleteMemoryUpdate::Object {
size,
removed_current_object,
}) => record_bucket_object_delete_memory(bucket, size, removed_current_object).await,
None => {}
}
}
/// `DeleteObjects` is idempotent. A raw filesystem `NotFound` can cross the /// `DeleteObjects` is idempotent. A raw filesystem `NotFound` can cross the
/// distributed delete path instead of its usual typed missing-object error. /// distributed delete path instead of its usual typed missing-object error.
fn is_delete_objects_not_found(error: &EcstoreError) -> bool { fn is_delete_objects_not_found(error: &EcstoreError) -> bool {
@@ -8409,8 +8458,6 @@ impl DefaultObjectUsecase {
object: ObjectToDelete, object: ObjectToDelete,
versioned: bool, versioned: bool,
version_suspended: bool, version_suspended: bool,
size: i64,
existing: Option<ObjectInfo>,
} }
// Phase 2 (bounded concurrency, backlog#929 / HP-8): collect the // Phase 2 (bounded concurrency, backlog#929 / HP-8): collect the
@@ -8428,32 +8475,23 @@ impl DefaultObjectUsecase {
skip_stat, skip_stat,
} = prepared; } = prepared;
let synthetic_version_id = object.version_id.is_none() && is_dir_object(&object.object_name); let synthetic_version_id = object.version_id.is_none() && is_dir_object(&object.object_name);
let (goi, source_missing) = if skip_stat { if !skip_stat {
(ObjectInfo::default(), false)
} else {
match store_ref.get_object_info(bucket_ref, &object.object_name, &opts).await { match store_ref.get_object_info(bucket_ref, &object.object_name, &opts).await {
Ok(res) => (res, false), Ok(_) => {}
Err(err) if is_err_object_not_found(&err) || is_err_version_not_found(&err) => { Err(err) if is_err_object_not_found(&err) || is_err_version_not_found(&err) => {}
(ObjectInfo::default(), true)
}
Err(err) => return Err(ApiError::from(err)), Err(err) => return Err(ApiError::from(err)),
} }
}; }
let size = goi.size;
if synthetic_version_id { if synthetic_version_id {
object.version_id = Some(Uuid::nil()); object.version_id = Some(Uuid::nil());
} }
let existing = (!skip_stat && !source_missing).then_some(goi);
Ok::<_, ApiError>(AdmittedDelete { Ok::<_, ApiError>(AdmittedDelete {
idx, idx,
object, object,
versioned: opts.versioned, versioned: opts.versioned,
version_suspended: opts.version_suspended, version_suspended: opts.version_suspended,
size,
existing,
}) })
})) }))
.buffered(DELETE_OBJECTS_PRE_STAT_CONCURRENCY) .buffered(DELETE_OBJECTS_PRE_STAT_CONCURRENCY)
@@ -8464,15 +8502,11 @@ impl DefaultObjectUsecase {
// per-key success/failure reporting is unchanged. // per-key success/failure reporting is unchanged.
let mut object_to_delete = Vec::new(); let mut object_to_delete = Vec::new();
let mut object_to_delete_idx = Vec::new(); let mut object_to_delete_idx = Vec::new();
let mut object_sizes = Vec::new();
let mut existing_object_infos = Vec::new();
let mut object_versioning = Vec::new(); let mut object_versioning = Vec::new();
for admitted in admitted_deletes { for admitted in admitted_deletes {
object_sizes.push(admitted.size);
object_to_delete_idx.push(admitted.idx); object_to_delete_idx.push(admitted.idx);
object_versioning.push((admitted.versioned, admitted.version_suspended)); object_versioning.push((admitted.versioned, admitted.version_suspended));
object_to_delete.push(admitted.object); object_to_delete.push(admitted.object);
existing_object_infos.push(admitted.existing);
} }
let cache_adapter = self.object_data_cache(); let cache_adapter = self.object_data_cache();
let cache_keys_before_delete = object_to_delete let cache_keys_before_delete = object_to_delete
@@ -8489,8 +8523,8 @@ impl DefaultObjectUsecase {
..Default::default() ..Default::default()
}; };
apply_bucket_generation_guard(&req, &bucket, &mut storage_delete_opts)?; apply_bucket_generation_guard(&req, &bucket, &mut storage_delete_opts)?;
let (dobjs, errs) = store let (dobjs, errs, accounting) = store
.delete_objects_with_tier_delete_journal(&bucket, object_to_delete.clone(), storage_delete_opts) .delete_objects_with_tier_delete_journal_and_accounting(&bucket, object_to_delete.clone(), storage_delete_opts)
.await; .await;
let _manager = get_concurrency_manager(); let _manager = get_concurrency_manager();
@@ -8515,17 +8549,16 @@ impl DefaultObjectUsecase {
delete_results[didx].delete_object = Some(deleted_object.clone()); delete_results[didx].delete_object = Some(deleted_object.clone());
let (versioned, version_suspended) = object_versioning[i]; let (versioned, version_suspended) = object_versioning[i];
let creates_delete_marker = object_to_delete[i].version_id.is_none() && versioned && !version_suspended; let creates_delete_marker = object_to_delete[i].version_id.is_none() && versioned && !version_suspended;
if creates_delete_marker { let committed_delete_marker = dobjs[i].delete_marker;
record_bucket_delete_marker_memory(&bucket).await; let delete_accounting = accounting.get(i).and_then(Option::as_ref);
} else { let update = delete_memory_update(
let size = object_sizes[i].max(0) as u64; creates_delete_marker,
record_bucket_object_delete_memory( committed_delete_marker,
&bucket, delete_request_targets_current(object_to_delete[i].version_id),
size, delete_accounting.and_then(|value| value.size),
existing_object_infos[i].is_some() && object_to_delete[i].version_id.is_none(), delete_accounting.is_some_and(|value| value.removed_current_object),
) );
.await; apply_delete_memory_update(&bucket, update).await;
}
} }
Err(error) => { Err(error) => {
delete_results[didx].error = Some(error); delete_results[didx].error = Some(error);
@@ -8803,12 +8836,24 @@ impl DefaultObjectUsecase {
let _ = invalidate_object_data_cache_after_delete_success(&cache_adapter, &bucket, &key).await; let _ = invalidate_object_data_cache_after_delete_success(&cache_adapter, &bucket, &key).await;
} }
// Fast in-memory update for immediate quota and admin usage consistency // Fast in-memory update for immediate quota and admin usage consistency.
if delete_creates_delete_marker(&opts) { // Prefix/force deletes and synthetic directory entries do not carry one
record_bucket_delete_marker_memory(&bucket).await; // committed object identity; leave their cache delta to reconciliation.
let update = if force_delete || obj_info.name.is_empty() || synthetic_version_id {
None
} else { } else {
record_bucket_object_delete_memory(&bucket, obj_info.size.max(0) as u64, opts.version_id.is_none()).await; // The storage commit returns this object's metadata while its
} // generation lock is held. Never fall back to a pre-delete stat:
// an overwrite can commit between that stat and this delete.
delete_memory_update(
delete_creates_delete_marker(&opts),
obj_info.delete_marker,
opts.version_id.is_none(),
quota_object_size(&obj_info).ok(),
delete_removes_current_object(&opts),
)
};
apply_delete_memory_update(&bucket, update).await;
if obj_info.name.is_empty() { if obj_info.name.is_empty() {
if let Some((operation_id, target_arns, generation)) = force_delete_intent { if let Some((operation_id, target_arns, generation)) = force_delete_intent {
@@ -17861,6 +17906,158 @@ mod tests {
assert!(!can_skip_delete_objects_pre_stat(false, &delete_marker_creating_opts(), false)); assert!(!can_skip_delete_objects_pre_stat(false, &delete_marker_creating_opts(), false));
} }
#[test]
fn delete_accounting_recognizes_explicit_null_as_current_object() {
let opts = ObjectOptions {
version_id: Some(Uuid::nil().to_string()),
version_suspended: true,
..Default::default()
};
assert!(delete_removes_current_object(&opts));
assert!(delete_request_targets_current(Some(Uuid::nil())));
assert!(!delete_request_targets_current(Some(Uuid::new_v4())));
assert!(!delete_removes_current_object(&ObjectOptions {
version_id: Some(Uuid::new_v4().to_string()),
..Default::default()
}));
}
#[test]
fn compressed_object_delete_restores_usage_baseline() {
let mut metadata = HashMap::new();
insert_str(&mut metadata, SUFFIX_COMPRESSION, "klauspost/compress/s2".to_string());
let object = ObjectInfo {
size: 400,
actual_size: 1000,
user_defined: Arc::new(metadata),
..Default::default()
};
let accounting_size = quota_object_size(&object).expect("logical compressed size should be canonical");
assert_eq!(
delete_memory_update(false, false, true, Some(accounting_size), true),
Some(DeleteMemoryUpdate::Object {
size: 1000,
removed_current_object: true,
})
);
}
#[test]
fn invalid_accounting_metadata_is_reconciled_without_overflow() {
assert_eq!(delete_memory_update(false, false, true, None, true), None);
assert_eq!(
delete_memory_update(false, true, true, None, true),
Some(DeleteMemoryUpdate::DeleteMarker)
);
}
#[tokio::test]
#[serial_test::serial]
async fn compressed_delete_requests_restore_usage_baseline() {
use crate::app::storage_api::test::contract::bucket::{BucketOperations as _, DeleteBucketOptions, MakeBucketOptions};
let store = crate::app::gating_test_env::shared_gating_ecstore().await;
if current_app_context().is_none() {
crate::app::runtime_sources::install_test_app_context(Arc::clone(&store)).await;
}
let bucket = format!("compressed-delete-request-{}", Uuid::new_v4().simple());
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create compressed delete request bucket");
// Seed the process-local usage with the canonical logical bytes. The
// direct storage PUT below intentionally does not apply an app-layer
// usage delta; the two real DELETE requests must remove exactly this
// amount through their request-layer wiring.
crate::app::storage_api::test::data_usage::seed_bucket_usage_memory_for_test(&bucket, 2_000).await;
for object in ["single", "batch"] {
let mut metadata = HashMap::new();
insert_str(&mut metadata, SUFFIX_COMPRESSION, "klauspost/compress/s2".to_string());
insert_str(&mut metadata, SUFFIX_ACTUAL_SIZE, "1000".to_string());
let reader = HashReader::from_stream(std::io::Cursor::new(vec![0x5a; 400]), 400, 1000, None, None, false)
.expect("compressed fixture reader should be valid");
let mut reader = PutObjReader::new(reader);
store
.put_object(
&bucket,
object,
&mut reader,
&ObjectOptions {
user_defined: metadata,
..Default::default()
},
)
.await
.expect("compressed fixture object should be written");
}
let mut single_req = build_request(
DeleteObjectInput::builder()
.bucket(bucket.clone())
.key("single".to_string())
.build()
.expect("single delete input should build"),
Method::DELETE,
);
single_req.extensions.insert(crate::storage::access::ReqInfo {
cred: Some(rustfs_credentials::Credentials::default()),
is_owner: true,
..Default::default()
});
DefaultObjectUsecase::from_global()
.execute_delete_object(single_req)
.await
.expect("single compressed delete should succeed");
assert_eq!(
crate::app::storage_api::test::data_usage::get_bucket_usage_memory(&bucket).await,
Some(1_000),
"single delete must subtract the logical accounting size"
);
let mut batch_req = build_request(
DeleteObjectsInput::builder()
.bucket(bucket.clone())
.delete(Delete {
objects: vec![ObjectIdentifier {
key: "batch".to_string(),
..Default::default()
}],
quiet: None,
})
.build()
.expect("batch delete input should build"),
Method::POST,
);
batch_req.extensions.insert(crate::storage::access::ReqInfo {
cred: Some(rustfs_credentials::Credentials::default()),
is_owner: true,
..Default::default()
});
DefaultObjectUsecase::from_global()
.execute_delete_objects(batch_req)
.await
.expect("batch compressed delete should succeed");
assert_eq!(
crate::app::storage_api::test::data_usage::get_bucket_usage_memory(&bucket).await,
Some(0),
"batch delete must subtract the committed logical accounting size"
);
store
.delete_bucket(
&bucket,
&DeleteBucketOptions {
force: true,
..Default::default()
},
)
.await
.expect("clean up compressed delete request bucket");
}
#[tokio::test] #[tokio::test]
async fn execute_get_object_attributes_returns_internal_error_when_store_uninitialized() { async fn execute_get_object_attributes_returns_internal_error_when_store_uninitialized() {
let input = GetObjectAttributesInput::builder() let input = GetObjectAttributesInput::builder()
+9 -1
View File
@@ -72,6 +72,11 @@ pub(crate) mod data_usage {
compute_bucket_usage, live_bucket_usage_computations, seed_bucket_usage_memory_for_test, store_data_usage_in_backend, compute_bucket_usage, live_bucket_usage_computations, seed_bucket_usage_memory_for_test, store_data_usage_in_backend,
}; };
#[cfg(test)]
pub(crate) async fn get_bucket_usage_memory(bucket: &str) -> Option<u64> {
crate::storage::storage_api::ecstore_data_usage::get_bucket_usage_memory(bucket).await
}
pub(crate) async fn record_bucket_object_delete_memory(bucket: &str, deleted_size: u64, removed_current_object: bool) { pub(crate) async fn record_bucket_object_delete_memory(bucket: &str, deleted_size: u64, removed_current_object: bool) {
crate::storage::storage_api::ecstore_data_usage::record_bucket_object_delete_memory( crate::storage::storage_api::ecstore_data_usage::record_bucket_object_delete_memory(
bucket, bucket,
@@ -1233,7 +1238,10 @@ pub(crate) mod test {
pub(crate) use super::access::ReqInfo; pub(crate) use super::access::ReqInfo;
pub(crate) use super::options::VERSIONING_CONFIG_LOOKUPS; pub(crate) use super::options::VERSIONING_CONFIG_LOOKUPS;
pub(crate) use super::{bucket, data_usage, ecfs, object_utils, runtime}; pub(crate) use super::{bucket, ecfs, object_utils, runtime};
pub(crate) mod data_usage {
pub(crate) use super::super::data_usage::*;
}
pub(crate) use crate::storage::storage_api::test_consumer::{get_global_bucket_metadata_sys, set_bucket_metadata}; pub(crate) use crate::storage::storage_api::test_consumer::{get_global_bucket_metadata_sys, set_bucket_metadata};
pub(crate) use crate::storage::storage_api::{ pub(crate) use crate::storage::storage_api::{
ECStore, Endpoint, Endpoints, PoolEndpoints, StorageObjectInfo, StorageObjectOptions, StoragePutObjReader, ECStore, Endpoint, Endpoints, PoolEndpoints, StorageObjectInfo, StorageObjectOptions, StoragePutObjReader,
+43 -1
View File
@@ -296,7 +296,10 @@ pub(crate) fn build_list_objects_v2_output(
let mut obj = Object { let mut obj = Object {
key: Some(key), key: Some(key),
last_modified: v.mod_time.map(Timestamp::from), last_modified: v.mod_time.map(Timestamp::from),
size: Some(v.get_actual_size().unwrap_or_default()), // Compressed legacy objects may retain an unknown (-1)
// logical-size sentinel; never expose that internal value in
// an S3 response.
size: Some(v.get_actual_size_or_physical()),
e_tag: v.etag.clone().map(|etag| to_s3s_etag(&etag)), e_tag: v.etag.clone().map(|etag| to_s3s_etag(&etag)),
storage_class: v.storage_class.clone().map(ObjectStorageClass::from), storage_class: v.storage_class.clone().map(ObjectStorageClass::from),
..Default::default() ..Default::default()
@@ -656,6 +659,45 @@ mod tests {
assert_eq!(output.common_prefixes.as_ref().map(std::vec::Vec::len), Some(2)); assert_eq!(output.common_prefixes.as_ref().map(std::vec::Vec::len), Some(2));
} }
#[test]
fn list_objects_never_exposes_compressed_unknown_size_sentinel() {
let mut metadata = std::collections::HashMap::new();
rustfs_utils::http::insert_str(
&mut metadata,
rustfs_utils::http::SUFFIX_COMPRESSION,
"klauspost/compress/s2".to_string(),
);
let output = build_list_objects_v2_output(
ListObjectsV2Info {
objects: vec![ObjectInfo {
name: "legacy-compressed".to_string(),
size: 128,
actual_size: -1,
user_defined: std::sync::Arc::new(metadata),
..Default::default()
}],
..Default::default()
},
false,
1000,
"bucket".to_string(),
String::new(),
None,
None,
None,
None,
);
assert_eq!(
output
.contents
.as_ref()
.and_then(|objects| objects.first())
.and_then(|object| object.size),
Some(128)
);
}
#[test] #[test]
fn list_responses_report_standard_for_legacy_label_only_file_metadata() { fn list_responses_report_standard_for_legacy_label_only_file_metadata() {
let version_id = Uuid::parse_str("11111111-2222-3333-4444-555555555555").expect("fixture version ID should be valid"); let version_id = Uuid::parse_str("11111111-2222-3333-4444-555555555555").expect("fixture version ID should be valid");
+2
View File
@@ -429,6 +429,8 @@ pub(crate) mod ecstore_config {
} }
pub(crate) mod ecstore_data_usage { pub(crate) mod ecstore_data_usage {
#[cfg(test)]
pub(crate) use rustfs_ecstore::api::data_usage::get_bucket_usage_memory;
pub(crate) use rustfs_ecstore::api::data_usage::{ pub(crate) use rustfs_ecstore::api::data_usage::{
apply_bucket_usage_memory_overlay, init_compression_total_memory_from_backend, load_admin_data_usage_from_backend_cached, apply_bucket_usage_memory_overlay, init_compression_total_memory_from_backend, load_admin_data_usage_from_backend_cached,
load_data_usage_from_backend, quota_object_size, record_bucket_delete_marker_memory, record_bucket_object_delete_memory, load_data_usage_from_backend, quota_object_size, record_bucket_delete_marker_memory, record_bucket_object_delete_memory,