Compare commits

..

1 Commits

Author SHA1 Message Date
overtrue 60fb7dc372 fix(ecstore): preserve remote delete error types 2026-08-22 17:05:52 +08:00
47 changed files with 302 additions and 4235 deletions
Generated
-1
View File
@@ -9587,7 +9587,6 @@ dependencies = [
"serde", "serde",
"serde_json", "serde_json",
"serial_test", "serial_test",
"sha2 0.11.0",
"temp-env", "temp-env",
"tempfile", "tempfile",
"thiserror 2.0.20", "thiserror 2.0.20",
+18 -64
View File
@@ -585,12 +585,9 @@ impl VersionsHistogram {
} }
} }
/// Replication statistics for a single target. /// Replication statistics for a single target
/// #[derive(Debug, Default, Clone, Serialize, Deserialize)]
/// Renamed from `ReplicationStats`; serde field names are preserved pub struct ReplicationStats {
/// byte-identically to maintain wire compatibility with existing snapshots.
#[derive(Debug, Default, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct ReplicationTargetUsage {
pub pending_size: u64, pub pending_size: u64,
pub replicated_size: u64, pub replicated_size: u64,
pub failed_size: u64, pub failed_size: u64,
@@ -603,7 +600,7 @@ pub struct ReplicationTargetUsage {
pub replicated_count: u64, pub replicated_count: u64,
} }
impl ReplicationTargetUsage { impl ReplicationStats {
pub fn is_empty(&self) -> bool { pub fn is_empty(&self) -> bool {
let Self { let Self {
pending_size, pending_size,
@@ -639,7 +636,7 @@ impl ReplicationTargetUsage {
/// Replication statistics for all targets /// Replication statistics for all targets
#[derive(Debug, Default, Clone, Serialize, Deserialize)] #[derive(Debug, Default, Clone, Serialize, Deserialize)]
pub struct ReplicationAllStats { pub struct ReplicationAllStats {
pub targets: HashMap<String, ReplicationTargetUsage>, pub targets: HashMap<String, ReplicationStats>,
pub replica_size: u64, pub replica_size: u64,
pub replica_count: u64, pub replica_count: u64,
} }
@@ -652,7 +649,7 @@ impl ReplicationAllStats {
targets, targets,
} = self; } = self;
*replica_size == 0 && *replica_count == 0 && targets.values().all(ReplicationTargetUsage::is_empty) *replica_size == 0 && *replica_count == 0 && targets.values().all(ReplicationStats::is_empty)
} }
#[deprecated(note = "use is_empty instead")] #[deprecated(note = "use is_empty instead")]
@@ -2469,7 +2466,7 @@ mod tests {
#[test] #[test]
fn replication_stats_empty_checks_every_field() { fn replication_stats_empty_checks_every_field() {
type SetField = fn(&mut ReplicationTargetUsage); type SetField = fn(&mut ReplicationStats);
let cases: [(&str, SetField); 10] = [ let cases: [(&str, SetField); 10] = [
("pending_size", |stats| stats.pending_size = 1), ("pending_size", |stats| stats.pending_size = 1),
@@ -2484,9 +2481,9 @@ mod tests {
("replicated_count", |stats| stats.replicated_count = 1), ("replicated_count", |stats| stats.replicated_count = 1),
]; ];
assert!(ReplicationTargetUsage::default().is_empty()); assert!(ReplicationStats::default().is_empty());
for (field, set_nonzero) in cases { for (field, set_nonzero) in cases {
let mut stats = ReplicationTargetUsage::default(); let mut stats = ReplicationStats::default();
set_nonzero(&mut stats); set_nonzero(&mut stats);
assert!(!stats.is_empty(), "{field} must make replication stats non-empty"); assert!(!stats.is_empty(), "{field} must make replication stats non-empty");
} }
@@ -2517,17 +2514,17 @@ mod tests {
} }
let empty_targets = ReplicationAllStats { let empty_targets = ReplicationAllStats {
targets: HashMap::from([("arn:test:empty".to_string(), ReplicationTargetUsage::default())]), targets: HashMap::from([("arn:test:empty".to_string(), ReplicationStats::default())]),
..Default::default() ..Default::default()
}; };
assert!(empty_targets.is_empty(), "all-empty targets must keep aggregate stats empty"); assert!(empty_targets.is_empty(), "all-empty targets must keep aggregate stats empty");
let stats = ReplicationAllStats { let stats = ReplicationAllStats {
targets: HashMap::from([ targets: HashMap::from([
("arn:test:empty".to_string(), ReplicationTargetUsage::default()), ("arn:test:empty".to_string(), ReplicationStats::default()),
( (
"arn:test:non-empty".to_string(), "arn:test:non-empty".to_string(),
ReplicationTargetUsage { ReplicationStats {
pending_count: 1, pending_count: 1,
..Default::default() ..Default::default()
}, },
@@ -2568,7 +2565,7 @@ mod tests {
replication_stats: Some(ReplicationAllStats { replication_stats: Some(ReplicationAllStats {
targets: HashMap::from([( targets: HashMap::from([(
"arn:test:pending".to_string(), "arn:test:pending".to_string(),
ReplicationTargetUsage { ReplicationStats {
pending_count: 1, pending_count: 1,
..Default::default() ..Default::default()
}, },
@@ -2717,7 +2714,7 @@ mod tests {
targets: HashMap::from([ targets: HashMap::from([
( (
"arn:self-only".to_string(), "arn:self-only".to_string(),
ReplicationTargetUsage { ReplicationStats {
pending_size: 7, pending_size: 7,
pending_count: 1, pending_count: 1,
..Default::default() ..Default::default()
@@ -2725,7 +2722,7 @@ mod tests {
), ),
( (
"arn:shared".to_string(), "arn:shared".to_string(),
ReplicationTargetUsage { ReplicationStats {
failed_size: 3, failed_size: 3,
failed_count: 1, failed_count: 1,
missed_threshold_size: 2, missed_threshold_size: 2,
@@ -2744,7 +2741,7 @@ mod tests {
targets: HashMap::from([ targets: HashMap::from([
( (
"arn:shared".to_string(), "arn:shared".to_string(),
ReplicationTargetUsage { ReplicationStats {
failed_size: 5, failed_size: 5,
failed_count: 2, failed_count: 2,
after_threshold_size: 4, after_threshold_size: 4,
@@ -2754,7 +2751,7 @@ mod tests {
), ),
( (
"arn:other-only".to_string(), "arn:other-only".to_string(),
ReplicationTargetUsage { ReplicationStats {
replicated_size: 11, replicated_size: 11,
replicated_count: 3, replicated_count: 3,
..Default::default() ..Default::default()
@@ -2996,9 +2993,7 @@ mod tests {
fn replication_target_deserialization_preserves_large_historical_maps() { fn replication_target_deserialization_preserves_large_historical_maps() {
let mut stats = ReplicationAllStats::default(); let mut stats = ReplicationAllStats::default();
for index in 0..=1024 { for index in 0..=1024 {
stats stats.targets.insert(format!("target-{index}"), ReplicationStats::default());
.targets
.insert(format!("target-{index}"), ReplicationTargetUsage::default());
} }
let encoded = rmp_serde::to_vec_named(&stats).expect("large replication target fixture should encode"); let encoded = rmp_serde::to_vec_named(&stats).expect("large replication target fixture should encode");
let decoded = rmp_serde::from_slice::<ReplicationAllStats>(&encoded) let decoded = rmp_serde::from_slice::<ReplicationAllStats>(&encoded)
@@ -3007,47 +3002,6 @@ mod tests {
assert_eq!(decoded.targets.len(), stats.targets.len()); assert_eq!(decoded.targets.len(), stats.targets.len());
} }
/// Round-trip test: encoding a [`ReplicationTargetUsage`] and decoding it back
/// must produce the exact same value. This guards against accidental serde
/// field-name drift during the `ReplicationStats` -> `ReplicationTargetUsage`
/// rename. Wire-level field names are the serialized Rust field identifiers,
/// which must remain byte-identical.
#[test]
fn replication_target_usage_rmp_round_trip() {
let original = ReplicationTargetUsage {
pending_size: 100,
replicated_size: 2_000,
failed_size: 50,
failed_count: 3,
pending_count: 7,
missed_threshold_size: 11,
after_threshold_size: 22,
missed_threshold_count: 1,
after_threshold_count: 2,
replicated_count: 99,
};
let buf = rmp_serde::to_vec_named(&original).expect("encode ReplicationTargetUsage to msgpack");
let decoded: ReplicationTargetUsage = rmp_serde::from_slice(&buf).expect("decode ReplicationTargetUsage from msgpack");
assert_eq!(original, decoded, "round-trip through rmp must preserve every field");
// Also verify that encoding as an unnamed sequence and then decoding
// with named fields produces the correct mapping (this catches reordering).
let named_buf = rmp_serde::to_vec_named(&original).expect("re-encode for field-name pinning");
// Spot-check that known field names appear in the named encoding.
let named_str = String::from_utf8_lossy(&named_buf);
assert!(named_str.contains("pending_size"), "field 'pending_size' must survive the rename");
assert!(named_str.contains("replicated_size"), "field 'replicated_size' must survive the rename");
assert!(
named_str.contains("missed_threshold_size"),
"field 'missed_threshold_size' must survive the rename"
);
assert!(
named_str.contains("after_threshold_count"),
"field 'after_threshold_count' must survive the rename"
);
}
#[test] #[test]
fn checked_merge_rejects_noncanonical_histograms_without_mutation() { fn checked_merge_rejects_noncanonical_histograms_without_mutation() {
let mut entry = DataUsageEntry { let mut entry = DataUsageEntry {
@@ -380,24 +380,10 @@ mod tests {
cluster.start_node(1).await?; cluster.start_node(1).await?;
let status_url = format!("{}/rustfs/admin/v3/background-heal/status", cluster.nodes[0].url); let status_url = format!("{}/rustfs/admin/v3/background-heal/status", cluster.nodes[0].url);
let mut recovered = serde_json::Value::Null; let status_body = signed_admin_post(&status_url, None, &cluster.access_key, &cluster.secret_key).await?;
for _ in 0..60 { assert!(
let status_body = signed_admin_post(&status_url, None, &cluster.access_key, &cluster.secret_key).await?; !status_body.contains("MissingContentLength"),
assert!( "background heal status should not fail without an explicit Content-Length: {status_body}"
!status_body.contains("MissingContentLength"),
"background heal status should not fail without an explicit Content-Length: {status_body}"
);
recovered = serde_json::from_str(&status_body)
.map_err(|err| format!("background heal status is not JSON ({err}): {status_body}"))?;
if recovered["clusterStatusComplete"] == serde_json::Value::Bool(true) {
break;
}
sleep(Duration::from_secs(1)).await;
}
assert_eq!(
recovered["clusterStatusComplete"],
serde_json::Value::Bool(true),
"cluster heal status should recover before root heal starts: {recovered}"
); );
let heal_body = r#"{"recursive":true,"dryRun":false,"remove":false,"recreate":true,"scanMode":2,"updateParity":false,"nolock":false}"#; let heal_body = r#"{"recursive":true,"dryRun":false,"remove":false,"recreate":true,"scanMode":2,"updateParity":false,"nolock":false}"#;
+2 -2
View File
@@ -317,6 +317,8 @@ 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,
@@ -328,8 +330,6 @@ 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 {
@@ -3425,44 +3425,6 @@ mod tests {
assert!(mutexes.contains_key("second")); assert!(mutexes.contains_key("second"));
} }
#[tokio::test]
async fn update_all_targets_publishes_disable_proxy_on_target_client() {
// The read-proxy selector (replication_proxy::get_proxy_targets) skips
// targets whose TargetClient carries disable_proxy — the persisted
// per-target opt-out must survive client publication.
let sys = BucketTargetSys::default();
let target = |arn: &str, disable_proxy: bool| BucketTarget {
arn: arn.to_string(),
endpoint: "192.168.1.10:9000".to_string(),
target_bucket: "target-bucket".to_string(),
region: "us-east-1".to_string(),
disable_proxy,
credentials: Some(Credentials {
access_key: "access".to_string(),
secret_key: "secret".to_string(),
session_token: None,
expiration: None,
}),
..Default::default()
};
let targets = BucketTargets {
targets: vec![target("arn:proxied", false), target("arn:opted-out", true)],
};
sys.update_all_targets("bucket", Some(&targets)).await;
let proxied = sys
.get_remote_target_client("bucket", "arn:proxied")
.await
.expect("client should be published");
assert!(!proxied.disable_proxy);
let opted_out = sys
.get_remote_target_client("bucket", "arn:opted-out")
.await
.expect("client should be published");
assert!(opted_out.disable_proxy, "disable_proxy must reach the published TargetClient");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)] #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn target_updates_serialize_client_build_through_publication_per_bucket() { async fn target_updates_serialize_client_build_through_publication_per_bucket() {
let sys = Arc::new(BucketTargetSys::default()); let sys = Arc::new(BucketTargetSys::default());
@@ -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_or_physical(); let asz = oi.get_actual_size().unwrap_or_default();
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_or_physical(); let actual_size = oi.get_actual_size().unwrap_or_default();
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_or_physical(), actual_size: object_info.get_actual_size().unwrap_or_default(),
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,20 +542,6 @@ 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 {
+1 -4
View File
@@ -248,13 +248,10 @@ impl Config {
let shard_size = shard_size as usize; let shard_size = shard_size as usize;
// Keep the historical two-data-shard object budget while preventing // Keep the historical two-data-shard object budget while preventing
// wider EC layouts from multiplying the maximum inline object size. // wider EC layouts from multiplying the maximum inline object size.
// Use div_ceil to match the shard_file_size calculation (which also uses
// div_ceil), avoiding a 1-byte rounding discrepancy that prevents inline
// for objects right at the threshold.
let inline_block = if self.initialized && self.inline_block_explicit { let inline_block = if self.initialized && self.inline_block_explicit {
self.inline_block self.inline_block
} else { } else {
DEFAULT_INLINE_OBJECT_BUDGET.div_ceil(data_shards).min(DEFAULT_INLINE_BLOCK) (DEFAULT_INLINE_OBJECT_BUDGET / data_shards).min(DEFAULT_INLINE_BLOCK)
}; };
if versioned { if versioned {
+60 -63
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::{DeleteAccounting, DeletedObject, ObjectIO as _, ObjectOperations as _, ObjectToDelete}, object::{DeletedObject, ObjectIO as _, ObjectOperations as _, ObjectToDelete},
range::HTTPRangeSpec, range::HTTPRangeSpec,
}; };
use crate::{ use crate::{
@@ -414,66 +414,6 @@ 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;
@@ -715,8 +655,65 @@ 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>>) {
let (deleted, errors, _) = self.delete_objects_with_accounting(bucket, objects, opts).await; // Default return value
(deleted, errors) let mut del_objects = vec![DeletedObject::default(); objects.len()];
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))]
+8 -108
View File
@@ -1391,37 +1391,7 @@ impl BucketUsageAccumulator {
} }
pub fn quota_object_size(object: &ObjectInfo) -> Result<u64, Error> { pub fn quota_object_size(object: &ObjectInfo) -> Result<u64, Error> {
// A compressed object may carry -1 while the transformed size is unknown let logical_size = u64::try_from(object.get_actual_size().map_err(Error::other)?).map_err(|_| Error::PartMissingOrCorrupt)?;
// (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 {
@@ -1429,8 +1399,12 @@ 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 == -1 { let actual_size = if part.actual_size < 0 {
0 if object.is_compressed() {
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)?
}; };
@@ -1438,7 +1412,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.unwrap_or(0).max(persisted_part_size)) Ok(logical_size.max(persisted_part_size))
} }
type UsageVersionPage = StorageListObjectVersionsInfo<ObjectInfo>; type UsageVersionPage = StorageListObjectVersionsInfo<ObjectInfo>;
@@ -3346,80 +3320,6 @@ 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() {
+2 -10
View File
@@ -7949,15 +7949,10 @@ impl DiskAPI for LocalDisk {
use std::io::Write as _; use std::io::Write as _;
let file_path = self.io_get_object_path(volume, path)?; let file_path = self.io_get_object_path(volume, path)?;
let lock_path = file_path.with_extension("rustfs-cas.lock");
let path = path.to_string(); let path = path.to_string();
let sync_metadata = effective_durability(volume).syncs_commit_metadata(); let sync_metadata = effective_durability(volume).syncs_commit_metadata();
return Ok(tokio::task::spawn_blocking(move || { return Ok(tokio::task::spawn_blocking(move || {
// A persistent directory lock bounds metadata growth. Removing
// per-target lock files can split flock ownership across inodes.
let lock_path = file_path
.parent()
.ok_or_else(|| std::io::Error::new(ErrorKind::InvalidInput, "conditional file has no parent"))?
.join(".rustfs-cas.lock");
let lock = std::fs::OpenOptions::new() let lock = std::fs::OpenOptions::new()
.create(true) .create(true)
.truncate(false) .truncate(false)
@@ -21846,10 +21841,7 @@ mod test {
let marker_path = disk let marker_path = disk
.get_object_path(RUSTFS_META_BUCKET, HEALING_MARKER_PATH) .get_object_path(RUSTFS_META_BUCKET, HEALING_MARKER_PATH)
.expect("marker path should resolve"); .expect("marker path should resolve");
let lock_path = marker_path let lock_path = marker_path.with_extension("rustfs-cas.lock");
.parent()
.expect("marker path should have a parent")
.join(".rustfs-cas.lock");
let lock = std::fs::OpenOptions::new() let lock = std::fs::OpenOptions::new()
.create(true) .create(true)
.truncate(false) .truncate(false)
+4 -34
View File
@@ -689,9 +689,6 @@ 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);
} }
@@ -703,25 +700,10 @@ 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);
} }
if self.actual_size == -1 && self.parts.is_empty() { let mut actual_size = 0;
return Ok(-1); self.parts.iter().for_each(|part| {
} 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)));
} }
@@ -736,18 +718,6 @@ 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;
+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::{DeleteAccounting, DeletedObject, HTTPPreconditions, ObjectIO as _, ObjectOperations as _, ObjectToDelete}, object::{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;
+6 -174
View File
@@ -45,7 +45,6 @@ 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;
@@ -2123,27 +2122,14 @@ impl SetDisks {
let erasure = Arc::new(erasure_from_file_info(&fi, false)?); let erasure = Arc::new(erasure_from_file_info(&fi, false)?);
let put_object_size = known_put_object_storage_size(data.size()); let put_object_size = known_put_object_storage_size(data.size());
let shard_file_size_raw = erasure.shard_file_size(put_object_size);
let is_inline_buffer = let is_inline_buffer =
storage_class_config.should_inline(shard_file_size_raw, erasure.data_shards, opts.versioned); storage_class_config.should_inline(erasure.shard_file_size(put_object_size), erasure.data_shards, opts.versioned);
let collect_stage_timing = rustfs_io_metrics::put_stage_metrics_enabled() || issue3031_diag_enabled(); let collect_stage_timing = rustfs_io_metrics::put_stage_metrics_enabled() || issue3031_diag_enabled();
let shard_file_size = shard_file_size_raw; let shard_file_size = erasure.shard_file_size(put_object_size);
let shard_size = erasure.shard_size(); let shard_size = erasure.shard_size();
let write_path = classify_put_write_path(is_inline_buffer, put_object_size, fi.erasure.block_size); let write_path = classify_put_write_path(is_inline_buffer, put_object_size, fi.erasure.block_size);
let direct_inline_commit = matches!(write_path, SmallWritePath::Inline); let direct_inline_commit = matches!(write_path, SmallWritePath::Inline);
{
use std::io::Write;
let msg = format!(
"INLINE_DEBUG: bucket={} obj={} size={} shard_fs={} ds={} bs={} inline={} direct={} path={} iblock={} ver={}\n",
bucket, object, put_object_size, shard_file_size_raw, erasure.data_shards, fi.erasure.block_size,
is_inline_buffer, direct_inline_commit, write_path.metric_label(), storage_class_config.inline_block(), opts.versioned
);
if let Ok(mut f) = std::fs::OpenOptions::new().create(true).append(true).open("/tmp/rustfs_inline_debug.log") {
let _ = f.write_all(msg.as_bytes());
}
let _ = std::io::stderr().write_all(msg.as_bytes());
}
rustfs_io_metrics::record_put_object_path(write_path.metric_label()); rustfs_io_metrics::record_put_object_path(write_path.metric_label());
let writer_setup_stage_start = collect_stage_timing.then(Instant::now); let writer_setup_stage_start = collect_stage_timing.then(Instant::now);
let (mut writers, errors) = if direct_inline_commit { let (mut writers, errors) = if direct_inline_commit {
@@ -5669,18 +5655,7 @@ 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()
@@ -5770,7 +5745,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, accounting); return (del_objects, del_errs);
} }
}, },
} }
@@ -5817,22 +5792,6 @@ 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
@@ -5961,12 +5920,7 @@ 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
@@ -6012,7 +5966,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
}); });
} }
} }
return (del_objects, del_errs, accounting); return (del_objects, del_errs);
} }
let mut persisted_journal_entries = Vec::with_capacity(journal_entries.len()); let mut persisted_journal_entries = Vec::with_capacity(journal_entries.len());
@@ -6250,16 +6204,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
} }
} }
// An accounting identity is actionable only when the delete result is (del_objects, del_errs)
// 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))]
@@ -6588,12 +6533,6 @@ 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)
@@ -7885,113 +7824,6 @@ 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::{
DeleteAccounting, DeletedObject, HTTPPreconditions, ObjectIO, ObjectLockDeleteOptions, ObjectLockRetentionOptions, DeletedObject, HTTPPreconditions, ObjectIO, ObjectLockDeleteOptions, ObjectLockRetentionOptions, ObjectOperations,
ObjectOperations, ObjectPreconditionError, ObjectPreconditionPart, ObjectPreconditionState, ObjectToDelete, ObjectPreconditionError, ObjectPreconditionPart, ObjectPreconditionState, ObjectToDelete,
}; };
pub(crate) trait EcstoreObjectIO: pub(crate) trait EcstoreObjectIO:
+2 -276
View File
@@ -297,16 +297,10 @@ impl ECStore {
#[cfg(test)] #[cfg(test)]
mod tests { mod tests {
use super::*; use super::*;
use crate::bucket::metadata_sys;
use crate::core::pools::{PoolDecommissionInfo, PoolStatus}; use crate::core::pools::{PoolDecommissionInfo, PoolStatus};
use crate::disk::{DeleteOptions, DiskOption, format::FormatV3, new_disk}; use crate::disk::{DiskOption, format::FormatV3, new_disk};
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints}; use crate::layout::endpoints::{Endpoints, PoolEndpoints};
use crate::runtime::instance::InstanceContext;
use crate::storage_api_contracts::bucket::{BucketOperations, MakeBucketOptions};
use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations};
use crate::store::init_format::{load_format_erasure, save_format_file}; use crate::store::init_format::{load_format_erasure, save_format_file};
use crate::store::init_local_disks_with_instance_ctx;
use tokio_util::sync::CancellationToken;
async fn minimal_heal_pool(pool_idx: usize) -> Arc<Sets> { async fn minimal_heal_pool(pool_idx: usize) -> Arc<Sets> {
let format = FormatV3::new(1, 1); let format = FormatV3::new(1, 1);
@@ -353,51 +347,6 @@ mod tests {
} }
} }
async fn multi_pool_heal_store() -> (tempfile::TempDir, Arc<ECStore>, CancellationToken) {
let temp_dir = tempfile::tempdir().expect("multi-pool heal test directory should be created");
let mut pool_endpoints = Vec::new();
for pool_index in 0..2 {
let mut endpoints = Vec::new();
for disk_index in 0..4 {
let disk_path = temp_dir.path().join(format!("pool{pool_index}-disk{disk_index}"));
tokio::fs::create_dir_all(&disk_path)
.await
.expect("multi-pool heal test disk should be created");
let mut endpoint = Endpoint::try_from(disk_path.to_str().expect("disk path should be utf8"))
.expect("test endpoint should parse");
endpoint.set_pool_index(pool_index);
endpoint.set_set_index(0);
endpoint.set_disk_index(disk_index);
endpoints.push(endpoint);
}
pool_endpoints.push(PoolEndpoints {
legacy: false,
set_count: 1,
drives_per_set: 4,
endpoints: Endpoints::from(endpoints),
cmd_line: format!("heal-owner-pool-{pool_index}"),
platform: "test".to_string(),
});
}
let endpoint_pools = EndpointServerPools::from(pool_endpoints);
let instance_ctx = Arc::new(InstanceContext::new());
init_local_disks_with_instance_ctx(&instance_ctx, endpoint_pools.clone())
.await
.expect("multi-pool local disks should initialize");
let shutdown = CancellationToken::new();
let store = ECStore::new_with_instance_ctx(
"127.0.0.1:0".parse().expect("test address should parse"),
endpoint_pools,
shutdown.clone(),
instance_ctx,
)
.await
.expect("multi-pool test store should initialize");
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
(temp_dir, store, shutdown)
}
#[tokio::test] #[tokio::test]
async fn heal_object_pool_scope_selects_only_requested_pool() { async fn heal_object_pool_scope_selects_only_requested_pool() {
let store = minimal_heal_store().await; let store = minimal_heal_store().await;
@@ -557,229 +506,6 @@ mod tests {
} }
} }
#[tokio::test]
#[serial_test::serial]
async fn unscoped_heal_object_suspended_owner_semantics() {
let (_temp_dir, store, shutdown) = multi_pool_heal_store().await;
let bucket = format!("heal-owner-{}", Uuid::new_v4().simple());
let active_object = "active-owner";
let suspended_only_object = "suspended-only";
let duplicate_object = "duplicate-owner";
let marker_object = "marker-owner";
let quorum_object = "quorum-owner";
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created in all pools");
let mut active_reader = PutObjReader::from_vec(b"active owner".to_vec());
store.pools[0]
.put_object(&bucket, active_object, &mut active_reader, &ObjectOptions::default())
.await
.expect("active owner object should be written");
let active_disks = store.pools[0].disk_set[0].disks.read().await.clone();
let missing_active_disk = active_disks[0].clone().expect("active disk should be online");
missing_active_disk
.delete(
&bucket,
active_object,
DeleteOptions {
recursive: true,
immediate: true,
..Default::default()
},
)
.await
.expect("active owner shard should be removed for repair");
assert!(
missing_active_disk.read_xl(&bucket, active_object, false).await.is_err(),
"the active owner fixture must start with one missing metadata copy"
);
let mut suspended_reader = PutObjReader::from_vec(b"suspended owner".to_vec());
store.pools[1]
.put_object(&bucket, suspended_only_object, &mut suspended_reader, &ObjectOptions::default())
.await
.expect("suspended owner object should be written");
for (pool_index, mod_time) in [1_i64, 2_i64].into_iter().enumerate() {
let mut duplicate_reader = PutObjReader::from_vec(format!("duplicate-pool-{pool_index}").into_bytes());
store.pools[pool_index]
.put_object(
&bucket,
duplicate_object,
&mut duplicate_reader,
&ObjectOptions {
mod_time: Some(OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(mod_time)),
..Default::default()
},
)
.await
.expect("duplicate owner object should be written");
}
let duplicate_missing_disk = store.pools[0].disk_set[0].disks.read().await[0]
.clone()
.expect("duplicate active owner disk should be online");
duplicate_missing_disk
.delete(
&bucket,
duplicate_object,
DeleteOptions {
recursive: true,
immediate: true,
..Default::default()
},
)
.await
.expect("duplicate active owner shard should be removed for repair");
let history_version = Uuid::new_v4();
let mut history_reader = PutObjReader::from_vec(b"marker history".to_vec());
store.pools[0]
.put_object(
&bucket,
marker_object,
&mut history_reader,
&ObjectOptions {
versioned: true,
version_id: Some(history_version.to_string()),
mod_time: Some(OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(1)),
..Default::default()
},
)
.await
.expect("versioned marker history should be written");
store.pools[0]
.delete_object(
&bucket,
marker_object,
ObjectOptions {
versioned: true,
mod_time: Some(OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(2)),
..Default::default()
},
)
.await
.expect("delete marker should be written");
let mut quorum_reader = PutObjReader::from_vec(b"quorum boundary".to_vec());
store.pools[0]
.put_object(&bucket, quorum_object, &mut quorum_reader, &ObjectOptions::default())
.await
.expect("quorum boundary object should be written");
{
let mut pool_meta = store.pool_meta.write().await;
let mut next = PoolMeta::new(&store.pools, &pool_meta);
next.pools[1].decommission = Some(PoolDecommissionInfo {
start_time: Some(OffsetDateTime::UNIX_EPOCH),
..Default::default()
});
*pool_meta = next;
}
let (_, duplicate_owner) = store
.get_latest_object_info_with_idx(&bucket, duplicate_object, &ObjectOptions::default())
.await
.expect("duplicate owner should resolve");
assert_eq!(duplicate_owner, 1, "latest duplicate must win when all pools are eligible");
let (_, active_duplicate_owner) = store
.get_latest_object_info_with_idx(
&bucket,
duplicate_object,
&ObjectOptions {
skip_decommissioned: true,
..Default::default()
},
)
.await
.expect("active duplicate owner should resolve");
assert_eq!(
active_duplicate_owner, 0,
"suspended duplicate must be excluded from active owner selection"
);
let (duplicate_result, duplicate_err) = store
.handle_heal_object(&bucket, duplicate_object, "", &HealOpts::default())
.await
.expect("duplicate owner heal should complete through the production path");
assert_eq!(duplicate_result.object, duplicate_object);
assert!(duplicate_err.is_none(), "active duplicate should be repaired: {duplicate_err:?}");
assert!(
duplicate_missing_disk.read_xl(&bucket, duplicate_object, false).await.is_ok(),
"production heal must repair the active duplicate owner rather than the suspended owner"
);
let (marker_info, marker_owner) = store
.get_latest_object_info_with_idx(
&bucket,
marker_object,
&ObjectOptions {
skip_decommissioned: true,
versioned: true,
..Default::default()
},
)
.await
.expect("latest delete marker should resolve");
assert_eq!(marker_owner, 0);
assert!(marker_info.delete_marker, "latest version must preserve delete-marker semantics");
let (active_result, active_err) = store
.handle_heal_object(&bucket, active_object, "", &HealOpts::default())
.await
.expect("unscoped active-owner heal should complete");
assert_eq!(active_result.object, active_object);
assert!(active_err.is_none(), "active owner must be selected even with a suspended pool");
assert!(
missing_active_disk.read_xl(&bucket, active_object, false).await.is_ok(),
"active owner heal must write the missing disk metadata: result={active_result:?}, err={active_err:?}"
);
assert!(
store.pools[1]
.get_object_info(&bucket, active_object, &ObjectOptions::default())
.await
.is_err(),
"the suspended pool must not be written for an active-owner object"
);
let (suspended_result, suspended_err) = store
.handle_heal_object(&bucket, suspended_only_object, "", &HealOpts::default())
.await
.expect("unscoped suspended-only heal should return a terminal result");
assert!(suspended_result.object.is_empty());
assert!(matches!(suspended_err, Some(Error::FileNotFound)));
assert!(
store.pools[1]
.get_object_info(&bucket, suspended_only_object, &ObjectOptions::default())
.await
.is_ok(),
"suspended-only data must remain untouched when unscoped heal reports absent"
);
let (_, explicit_err) = store
.handle_heal_object(
&bucket,
suspended_only_object,
"",
&HealOpts {
pool: Some(1),
..Default::default()
},
)
.await
.expect("explicit suspended-owner heal should return a mapped error");
assert!(matches!(explicit_err, Some(Error::SlowDown)));
let original_quorum_disks = store.pools[0].disk_set[0].disks.read().await.clone();
let surviving_quorum_disk = original_quorum_disks[3].clone();
*store.pools[0].disk_set[0].disks.write().await = vec![None, None, None, surviving_quorum_disk];
let (_, quorum_err) = store
.handle_heal_object(&bucket, quorum_object, "", &HealOpts::default())
.await
.expect("quorum boundary heal should return a mapped result");
*store.pools[0].disk_set[0].disks.write().await = original_quorum_disks;
assert!(
matches!(quorum_err, Some(Error::ErasureReadQuorum)),
"quorum-boundary heal must preserve quorum error, got {quorum_err:?}"
);
shutdown.cancel();
}
#[tokio::test] #[tokio::test]
async fn handle_heal_format_continues_after_a_pool_error() { async fn handle_heal_format_continues_after_a_pool_error() {
let canonical_format = FormatV3::new(1, 3); let canonical_format = FormatV3::new(1, 3);
+11 -54
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::{DeleteAccounting, ObjectIO as _, ObjectOperations as _}, object::{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,14 +1216,6 @@ 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();
@@ -2320,22 +2312,6 @@ 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
@@ -2713,19 +2689,6 @@ 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()
@@ -2738,7 +2701,6 @@ 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() {
@@ -2752,7 +2714,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_with_accounting(objects.as_slice(), err), Err(err) => return return_batch_delete_lock_error(objects.as_slice(), err),
} }
}; };
if let Some(guard) = _bucket_lifecycle_guard.as_ref() { if let Some(guard) = _bucket_lifecycle_guard.as_ref() {
@@ -2764,21 +2726,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, accounting); return (del_objects, errors);
} }
} }
} }
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_with_accounting(objects.as_slice(), err); return return_batch_delete_lock_error(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_with_accounting(objects.as_slice(), err), Err(err) => return return_batch_delete_lock_error(objects.as_slice(), err),
}) })
}; };
if let Some(guard) = _object_lock_metadata_guard.as_ref() { if let Some(guard) = _object_lock_metadata_guard.as_ref() {
@@ -2788,7 +2750,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_with_accounting(objects.as_slice(), err), Err(err) => return return_batch_delete_lock_error(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,
@@ -2804,10 +2766,7 @@ 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_with_accounting( return return_batch_delete_lock_error(objects.as_slice(), StorageError::BucketNotFound(bucket.to_string()));
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() {
@@ -2815,7 +2774,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_with_accounting(objects.as_slice(), err), Err(err) => return return_batch_delete_lock_error(objects.as_slice(), err),
}; };
let mut futures = Vec::with_capacity(self.pools.len()); let mut futures = Vec::with_capacity(self.pools.len());
@@ -2824,24 +2783,22 @@ 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_with_accounting(bucket, objects.clone(), opts.clone())); futures.push(pool.delete_objects(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, pool_accounting) in results.iter() { for (dels, errs) 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();
} }
} }
} }
@@ -2850,7 +2807,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, accounting) (del_objects, del_errs)
// let mut futures = Vec::with_capacity(objects.len()); // let mut futures = Vec::with_capacity(objects.len());
-1
View File
@@ -91,7 +91,6 @@ metrics = { workspace = true }
base64 = { workspace = true } base64 = { workspace = true }
bytes = { workspace = true } bytes = { workspace = true }
crc-fast = { workspace = true } crc-fast = { workspace = true }
sha2 = { workspace = true }
[dev-dependencies] [dev-dependencies]
serde_json = { workspace = true, features = ["raw_value"] } serde_json = { workspace = true, features = ["raw_value"] }
-5
View File
@@ -373,11 +373,6 @@ impl ErasureSetHealer {
set_disk_id: &str, set_disk_id: &str,
buckets: &[String], buckets: &[String],
) -> Result<(ResumeManager, CheckpointManager)> { ) -> Result<(ResumeManager, CheckpointManager)> {
if self.replacement_task_id.is_none() && CheckpointManager::is_blocked(&self.disk, task_id).await {
return Err(Error::TaskExecutionFailed {
message: format!("Resume task {task_id} has a blocked checkpoint"),
});
}
// check if resume state exists // check if resume state exists
let has_resume_state = if self.replacement_task_id.is_some() { let has_resume_state = if self.replacement_task_id.is_some() {
ResumeManager::has_replacement_intent(&self.disk, task_id).await ResumeManager::has_replacement_intent(&self.disk, task_id).await
-1
View File
@@ -51,7 +51,6 @@ const RESUME_STATE_FILE: &str = "ahm_resume_state.json";
const REPLACEMENT_INTENT_FILE: &str = "ahm_replacement_intent.json"; const REPLACEMENT_INTENT_FILE: &str = "ahm_replacement_intent.json";
const RESUME_PROGRESS_FILE: &str = "ahm_progress.json"; const RESUME_PROGRESS_FILE: &str = "ahm_progress.json";
pub(super) const RESUME_CHECKPOINT_FILE: &str = "ahm_checkpoint.json"; pub(super) const RESUME_CHECKPOINT_FILE: &str = "ahm_checkpoint.json";
pub(super) const RESUME_CHECKPOINT_BLOCKED_FILE: &str = "ahm_checkpoint.blocked";
const REPLACEMENT_COMPLETION_PROOF_FILE: &str = "ahm_replacement_completion_proof.json"; const REPLACEMENT_COMPLETION_PROOF_FILE: &str = "ahm_replacement_completion_proof.json";
const REPLACEMENT_RECOVERY_DIR: &str = "ahm-replacement"; const REPLACEMENT_RECOVERY_DIR: &str = "ahm-replacement";
const REPLACEMENT_INTENT_SEAL_FILE: &str = "ahm_replacement_intent_seal"; const REPLACEMENT_INTENT_SEAL_FILE: &str = "ahm_replacement_intent_seal";
+21 -315
View File
@@ -13,31 +13,26 @@
// limitations under the License. // limitations under the License.
use crate::{Error, Result}; use crate::{Error, Result};
use base64::Engine as _;
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use std::collections::HashSet; use std::collections::HashSet;
use std::path::Path; use std::path::Path;
use std::sync::{Arc, Mutex}; use std::sync::{Arc, Mutex};
use std::time::{SystemTime, UNIX_EPOCH}; use std::time::{SystemTime, UNIX_EPOCH};
use tokio::sync::{Mutex as AsyncMutex, RwLock}; use tokio::sync::RwLock;
use tracing::{debug, warn}; use tracing::{debug, warn};
use super::super::storage_api::owner::{EcstoreConditionalFileUpdate, EcstoreDiskAPI, EcstoreDiskBytes}; use super::super::{BUCKET_META_PREFIX, DiskStore, HealDiskExt as _, RUSTFS_META_BUCKET};
use super::super::{BUCKET_META_PREFIX, DiskStore, HealDiskExt, RUSTFS_META_BUCKET};
use super::{ use super::{
LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, PersistThrottle, RESUME_CHECKPOINT_BLOCKED_FILE, RESUME_CHECKPOINT_FILE, LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, PersistThrottle, RESUME_CHECKPOINT_FILE, delete_resume_file, path_to_str,
delete_resume_file, path_to_str, validate_resume_task_id, validate_resume_task_id,
}; };
const EVENT_HEAL_CHECKPOINT_STATE: &str = "heal_checkpoint_state"; const EVENT_HEAL_CHECKPOINT_STATE: &str = "heal_checkpoint_state";
const RESUME_CHECKPOINT_DIGEST_FILE: &str = "ahm_checkpoint.sha256";
const CHECKPOINT_PER_VERSION_SCHEMA: u32 = 5;
/// Current on-disk schema version for `ResumeCheckpoint`. Same rationale as /// Current on-disk schema version for `ResumeCheckpoint`. Same rationale as
/// `CURRENT_RESUME_SCHEMA`: pre-per-version dedup identities are not comparable /// `CURRENT_RESUME_SCHEMA`: pre-per-version dedup identities are not comparable
/// to the new `compose_key` identities, so a stale checkpoint is discarded. /// to the new `compose_key` identities, so a stale checkpoint is discarded.
pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 6; pub(super) const CURRENT_CHECKPOINT_SCHEMA: u32 = 5;
/// resume checkpoint /// resume checkpoint
#[derive(Debug, Clone, Serialize, Deserialize)] #[derive(Debug, Clone, Serialize, Deserialize)]
@@ -62,11 +57,6 @@ pub struct ResumeCheckpoint {
pub failed_objects: HashSet<String>, pub failed_objects: HashSet<String>,
/// skipped objects /// skipped objects
pub skipped_objects: HashSet<String>, pub skipped_objects: HashSet<String>,
/// Integrity digest over the checkpoint with this field set to `None`.
/// Keeping it in the checkpoint makes the payload and its authentication
/// record one CAS generation instead of two independently-written files.
#[serde(default)]
pub integrity_digest: Option<String>,
} }
impl ResumeCheckpoint { impl ResumeCheckpoint {
@@ -80,7 +70,6 @@ impl ResumeCheckpoint {
processed_objects: HashSet::new(), processed_objects: HashSet::new(),
failed_objects: HashSet::new(), failed_objects: HashSet::new(),
skipped_objects: HashSet::new(), skipped_objects: HashSet::new(),
integrity_digest: None,
} }
} }
@@ -127,111 +116,17 @@ pub struct CheckpointManager {
disk: DiskStore, disk: DiskStore,
checkpoint: Arc<RwLock<ResumeCheckpoint>>, checkpoint: Arc<RwLock<ResumeCheckpoint>>,
throttle: Mutex<PersistThrottle>, throttle: Mutex<PersistThrottle>,
save_lock: AsyncMutex<()>,
last_saved: Mutex<Option<EcstoreDiskBytes>>,
} }
impl CheckpointManager { impl CheckpointManager {
fn blocked_path(task_id: &str) -> std::path::PathBuf {
Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}"))
}
/// Return whether a checkpoint was permanently isolated after a malformed
/// or unsupported snapshot was observed.
pub(crate) async fn is_blocked(disk: &DiskStore, task_id: &str) -> bool {
if validate_resume_task_id(task_id).is_err() {
return false;
}
let blocked_path = Self::blocked_path(task_id);
let Ok(path) = path_to_str(&blocked_path) else {
return false;
};
match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path).await {
Ok(_) => true,
Err(crate::heal::DiskError::FileNotFound) => false,
Err(_) => true,
}
}
/// Validate the checkpoint while enumerating resumable state. This reads
/// the checkpoint once and also isolates malformed or unsupported data.
pub(crate) async fn is_resumable(disk: &DiskStore, task_id: &str) -> Result<bool> {
validate_resume_task_id(task_id)?;
if Self::is_blocked(disk, task_id).await {
return Err(Error::InvalidCheckpoint(format!("Resume task {task_id} has a blocked checkpoint")));
}
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
let Ok(path) = path_to_str(&file_path) else {
return Err(Error::InvalidCheckpoint("Resume checkpoint path is not valid UTF-8".to_string()));
};
match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path).await {
Ok(bytes) if bytes.is_empty() => Ok(true),
Ok(bytes) => Self::load_from_data(disk.clone(), task_id, bytes.to_vec())
.await
.map(|_| true),
Err(crate::heal::DiskError::FileNotFound) => Ok(true),
Err(error) => Err(error.into()),
}
}
async fn block_invalid_snapshot(disk: &DiskStore, task_id: &str) {
// This marker is intentionally version-agnostic: an unsupported reader
// must stop selector retries until an operator cleans up the snapshot.
let blocked_path = Self::blocked_path(task_id);
let Ok(path) = path_to_str(&blocked_path) else {
return;
};
let result = EcstoreDiskAPI::compare_and_update_file(
disk.as_ref(),
RUSTFS_META_BUCKET,
path,
None,
Some(EcstoreDiskBytes::from_static(b"blocked")),
)
.await;
match result {
Ok(EcstoreConditionalFileUpdate::Updated | EcstoreConditionalFileUpdate::Mismatch) => {}
Ok(EcstoreConditionalFileUpdate::Missing) => warn!(
target: "rustfs::heal::resume",
event = EVENT_HEAL_CHECKPOINT_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_RESUME,
task_id,
state = "blocked_marker_write_failed",
error = "marker target disappeared",
"Heal checkpoint could not persist its blocked marker"
),
Err(error) => warn!(
target: "rustfs::heal::resume",
event = EVENT_HEAL_CHECKPOINT_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_RESUME,
task_id,
state = "blocked_marker_write_failed",
error = %error,
"Heal checkpoint could not persist its blocked marker"
),
}
}
/// create new checkpoint manager /// create new checkpoint manager
pub async fn new(disk: DiskStore, task_id: String) -> Result<Self> { pub async fn new(disk: DiskStore, task_id: String) -> Result<Self> {
validate_resume_task_id(&task_id)?; validate_resume_task_id(&task_id)?;
let checkpoint_volume = format!("{RUSTFS_META_BUCKET}/{BUCKET_META_PREFIX}");
if let Err(error) = EcstoreDiskAPI::make_volume(disk.as_ref(), &checkpoint_volume).await
&& error != crate::heal::DiskError::VolumeExists
{
return Err(Error::TaskExecutionFailed {
message: format!("Failed to create checkpoint volume: {error}"),
});
}
let checkpoint = ResumeCheckpoint::new(task_id); let checkpoint = ResumeCheckpoint::new(task_id);
let manager = Self { let manager = Self {
disk, disk,
checkpoint: Arc::new(RwLock::new(checkpoint)), checkpoint: Arc::new(RwLock::new(checkpoint)),
throttle: Mutex::new(PersistThrottle::new()), throttle: Mutex::new(PersistThrottle::new()),
save_lock: AsyncMutex::new(()),
last_saved: Mutex::new(None),
}; };
// save initial checkpoint // save initial checkpoint
@@ -245,7 +140,6 @@ impl CheckpointManager {
error = %e, error = %e,
"Heal checkpoint persistence failed" "Heal checkpoint persistence failed"
); );
return Err(e);
} }
Ok(manager) Ok(manager)
} }
@@ -254,22 +148,11 @@ impl CheckpointManager {
pub async fn load_from_disk(disk: DiskStore, task_id: &str) -> Result<Self> { pub async fn load_from_disk(disk: DiskStore, task_id: &str) -> Result<Self> {
validate_resume_task_id(task_id)?; validate_resume_task_id(task_id)?;
let checkpoint_data = Self::read_checkpoint_file(&disk, task_id).await?; let checkpoint_data = Self::read_checkpoint_file(&disk, task_id).await?;
Self::load_from_data(disk, task_id, checkpoint_data).await let mut checkpoint: ResumeCheckpoint =
} serde_json::from_slice(&checkpoint_data).map_err(|e| Error::TaskExecutionFailed {
message: format!("Failed to deserialize checkpoint: {e}"),
async fn load_from_data(disk: DiskStore, task_id: &str, checkpoint_data: Vec<u8>) -> Result<Self> { })?;
validate_resume_task_id(task_id)?;
let mut checkpoint: ResumeCheckpoint = match serde_json::from_slice(&checkpoint_data) {
Ok(checkpoint) => checkpoint,
Err(error) => {
Self::block_invalid_snapshot(&disk, task_id).await;
return Err(Error::TaskExecutionFailed {
message: format!("Failed to deserialize checkpoint: {error}"),
});
}
};
if checkpoint.task_id != task_id { if checkpoint.task_id != task_id {
Self::block_invalid_snapshot(&disk, task_id).await;
return Err(Error::TaskExecutionFailed { return Err(Error::TaskExecutionFailed {
message: "Resume checkpoint task id does not match filename".to_string(), message: "Resume checkpoint task id does not match filename".to_string(),
}); });
@@ -280,7 +163,6 @@ impl CheckpointManager {
// identities. Discard the stale sets and position, then stamp the // identities. Discard the stale sets and position, then stamp the
// current schema so the scan restarts cleanly. // current schema so the scan restarts cleanly.
if checkpoint.schema_version > CURRENT_CHECKPOINT_SCHEMA { if checkpoint.schema_version > CURRENT_CHECKPOINT_SCHEMA {
Self::block_invalid_snapshot(&disk, task_id).await;
return Err(Error::TaskExecutionFailed { return Err(Error::TaskExecutionFailed {
message: format!( message: format!(
"Checkpoint schema {} is newer than supported schema {CURRENT_CHECKPOINT_SCHEMA}", "Checkpoint schema {} is newer than supported schema {CURRENT_CHECKPOINT_SCHEMA}",
@@ -288,43 +170,7 @@ impl CheckpointManager {
), ),
}); });
} }
if checkpoint.schema_version < CURRENT_CHECKPOINT_SCHEMA {
if let Some(expected) = checkpoint.integrity_digest.as_deref() {
let actual = Self::checkpoint_digest(&Self::serialize_without_digest(&checkpoint)?);
if expected != actual {
Self::block_invalid_snapshot(&disk, task_id).await;
return Err(Error::InvalidCheckpoint(format!(
"Resume checkpoint digest does not match task {task_id}"
)));
}
} else if checkpoint.schema_version >= CURRENT_CHECKPOINT_SCHEMA {
Self::block_invalid_snapshot(&disk, task_id).await;
return Err(Error::InvalidCheckpoint(format!(
"Resume checkpoint digest is missing for task {task_id}"
)));
} else {
let digest_path = Self::digest_path(task_id);
let digest_path = path_to_str(&digest_path)?;
match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, digest_path).await {
Ok(expected) => {
let actual = Self::checkpoint_digest(&checkpoint_data);
if expected.as_ref() != actual.as_bytes() {
Self::block_invalid_snapshot(&disk, task_id).await;
return Err(Error::InvalidCheckpoint(format!(
"Resume checkpoint digest does not match task {task_id}"
)));
}
}
Err(crate::heal::DiskError::FileNotFound) => {}
Err(error) => {
return Err(Error::TaskExecutionFailed {
message: format!("Failed to read checkpoint digest: {error}"),
});
}
}
}
if checkpoint.schema_version < CHECKPOINT_PER_VERSION_SCHEMA {
warn!( warn!(
target: "rustfs::heal::resume", target: "rustfs::heal::resume",
event = EVENT_HEAL_CHECKPOINT_STATE, event = EVENT_HEAL_CHECKPOINT_STATE,
@@ -341,15 +187,13 @@ impl CheckpointManager {
checkpoint.skipped_objects.clear(); checkpoint.skipped_objects.clear();
checkpoint.current_bucket_index = 0; checkpoint.current_bucket_index = 0;
checkpoint.current_object_index = 0; checkpoint.current_object_index = 0;
checkpoint.schema_version = CURRENT_CHECKPOINT_SCHEMA;
} }
checkpoint.schema_version = CURRENT_CHECKPOINT_SCHEMA;
Ok(Self { Ok(Self {
disk, disk,
checkpoint: Arc::new(RwLock::new(checkpoint)), checkpoint: Arc::new(RwLock::new(checkpoint)),
throttle: Mutex::new(PersistThrottle::new()), throttle: Mutex::new(PersistThrottle::new()),
save_lock: AsyncMutex::new(()),
last_saved: Mutex::new(Some(EcstoreDiskBytes::from(checkpoint_data))),
}) })
} }
@@ -360,7 +204,7 @@ impl CheckpointManager {
} }
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}")); let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
match path_to_str(&file_path) { match path_to_str(&file_path) {
Ok(path_str) => match HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path_str).await { Ok(path_str) => match disk.read_all(RUSTFS_META_BUCKET, path_str).await {
Ok(data) => !data.is_empty(), Ok(data) => !data.is_empty(),
Err(_) => false, Err(_) => false,
}, },
@@ -448,8 +292,6 @@ impl CheckpointManager {
let checkpoint_file = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}")); let checkpoint_file = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
delete_resume_file(&self.disk, &checkpoint_file).await?; delete_resume_file(&self.disk, &checkpoint_file).await?;
delete_resume_file(&self.disk, &Self::digest_path(&task_id)).await?;
delete_resume_file(&self.disk, &Self::blocked_path(&task_id)).await?;
debug!( debug!(
target: "rustfs::heal::resume", target: "rustfs::heal::resume",
@@ -465,130 +307,21 @@ impl CheckpointManager {
/// save checkpoint to disk /// save checkpoint to disk
async fn save_checkpoint(&self) -> Result<()> { async fn save_checkpoint(&self) -> Result<()> {
// Serialize saves and take the snapshot only after acquiring the lock: let checkpoint = self.checkpoint.read().await;
// a slower writer must not publish a snapshot taken before a newer one.
let _save_guard = self.save_lock.lock().await;
let checkpoint = self.checkpoint.read().await.clone();
validate_resume_task_id(&checkpoint.task_id)?; validate_resume_task_id(&checkpoint.task_id)?;
let unsigned_checkpoint_data = Self::serialize_without_digest(&checkpoint)?; let checkpoint_data = serde_json::to_vec(&*checkpoint).map_err(|e| Error::TaskExecutionFailed {
let digest = Self::checkpoint_digest(&unsigned_checkpoint_data); message: format!("Failed to serialize checkpoint: {e}"),
let mut persisted_checkpoint = checkpoint.clone(); })?;
persisted_checkpoint.integrity_digest = Some(digest);
let checkpoint_data =
EcstoreDiskBytes::from(serde_json::to_vec(&persisted_checkpoint).map_err(|e| Error::TaskExecutionFailed {
message: format!("Failed to serialize checkpoint: {e}"),
})?);
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{}_{}", checkpoint.task_id, RESUME_CHECKPOINT_FILE)); let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{}_{}", checkpoint.task_id, RESUME_CHECKPOINT_FILE));
let path_str = path_to_str(&file_path)?; let path_str = path_to_str(&file_path)?;
let last_saved = self self.disk
.last_saved .write_all(RUSTFS_META_BUCKET, path_str, checkpoint_data.into())
.lock()
.map_err(|_| Error::TaskExecutionFailed {
message: "Checkpoint save state lock is poisoned; refusing to save".to_string(),
})?
.clone();
let update = EcstoreDiskAPI::compare_and_update_file(
self.disk.as_ref(),
RUSTFS_META_BUCKET,
path_str,
last_saved.clone(),
Some(checkpoint_data.clone()),
)
.await
.map_err(|e| Error::TaskExecutionFailed {
message: format!("Failed to save checkpoint: {e}"),
})?;
let expected = match update {
EcstoreConditionalFileUpdate::Updated => None,
EcstoreConditionalFileUpdate::Missing => {
return Err(Error::TaskExecutionFailed {
message: "Checkpoint was removed after this manager saved it; refusing to recreate it".to_string(),
});
}
EcstoreConditionalFileUpdate::Mismatch => {
// A healthy manager normally completes the CAS above without
// another read or JSON parse. Inspect only after a mismatch so
// corruption and future schemas cannot be overwritten blindly.
let existing = match HealDiskExt::read_all(self.disk.as_ref(), RUSTFS_META_BUCKET, path_str).await {
Ok(existing) => existing,
Err(crate::heal::DiskError::FileNotFound) => {
return Err(Error::TaskExecutionFailed {
message: "Checkpoint was removed after this manager saved it; refusing to recreate it".to_string(),
});
}
Err(error) => {
return Err(Error::TaskExecutionFailed {
message: format!("Failed to inspect checkpoint after CAS mismatch: {error}"),
});
}
};
if existing.is_empty() && last_saved.is_none() {
Some(existing)
} else {
let current: ResumeCheckpoint = match serde_json::from_slice(&existing) {
Ok(current) => current,
Err(error) => {
Self::block_invalid_snapshot(&self.disk, &checkpoint.task_id).await;
return Err(Error::TaskExecutionFailed {
message: format!("Existing checkpoint is corrupt: {error}"),
});
}
};
if current.task_id != checkpoint.task_id {
Self::block_invalid_snapshot(&self.disk, &checkpoint.task_id).await;
return Err(Error::TaskExecutionFailed {
message: "Existing checkpoint task id does not match filename".to_string(),
});
}
if current.schema_version > CURRENT_CHECKPOINT_SCHEMA {
Self::block_invalid_snapshot(&self.disk, &checkpoint.task_id).await;
return Err(Error::TaskExecutionFailed {
message: format!(
"Existing checkpoint schema {} is newer than supported schema {CURRENT_CHECKPOINT_SCHEMA}",
current.schema_version
),
});
}
if last_saved.as_ref().is_none_or(|saved| saved.as_ref() != existing.as_ref()) {
return Err(Error::TaskExecutionFailed {
message: "Checkpoint changed since this manager loaded it; refusing to overwrite newer progress"
.to_string(),
});
}
Some(existing)
}
}
};
if let Some(expected) = expected {
match EcstoreDiskAPI::compare_and_update_file(
self.disk.as_ref(),
RUSTFS_META_BUCKET,
path_str,
Some(expected),
Some(checkpoint_data.clone()),
)
.await .await
.map_err(|e| Error::TaskExecutionFailed { .map_err(|e| Error::TaskExecutionFailed {
message: format!("Failed to save checkpoint after CAS mismatch: {e}"), message: format!("Failed to save checkpoint: {e}"),
})? { })?;
EcstoreConditionalFileUpdate::Updated => {}
EcstoreConditionalFileUpdate::Missing | EcstoreConditionalFileUpdate::Mismatch => {
return Err(Error::TaskExecutionFailed {
message: "Checkpoint changed while saving; refusing to overwrite newer progress".to_string(),
});
}
}
}
let mut last_saved = self.last_saved.lock().map_err(|_| Error::TaskExecutionFailed {
message: "Checkpoint save state lock is poisoned after save".to_string(),
})?;
*last_saved = Some(checkpoint_data);
debug!( debug!(
target: "rustfs::heal::resume", target: "rustfs::heal::resume",
@@ -608,38 +341,11 @@ impl CheckpointManager {
let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}")); let file_path = Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"));
let path_str = path_to_str(&file_path)?; let path_str = path_to_str(&file_path)?;
HealDiskExt::read_all(disk.as_ref(), RUSTFS_META_BUCKET, path_str) disk.read_all(RUSTFS_META_BUCKET, path_str)
.await .await
.map(|bytes| bytes.to_vec()) .map(|bytes| bytes.to_vec())
.map_err(|e| Error::TaskExecutionFailed { .map_err(|e| Error::TaskExecutionFailed {
message: format!("Failed to read checkpoint file: {e}"), message: format!("Failed to read checkpoint file: {e}"),
}) })
} }
fn serialize_without_digest(checkpoint: &ResumeCheckpoint) -> Result<Vec<u8>> {
let mut unsigned = checkpoint.clone();
unsigned.integrity_digest = None;
let mut value = serde_json::to_value(&unsigned).map_err(|e| Error::TaskExecutionFailed {
message: format!("Failed to serialize checkpoint: {e}"),
})?;
for field in ["processed_objects", "failed_objects", "skipped_objects"] {
let Some(values) = value.get_mut(field).and_then(serde_json::Value::as_array_mut) else {
return Err(Error::TaskExecutionFailed {
message: format!("Failed to canonicalize checkpoint field: {field}"),
});
};
values.sort_by(|left, right| left.as_str().cmp(&right.as_str()));
}
serde_json::to_vec(&value).map_err(|e| Error::TaskExecutionFailed {
message: format!("Failed to serialize checkpoint: {e}"),
})
}
fn checkpoint_digest(checkpoint_data: &[u8]) -> String {
base64::engine::general_purpose::STANDARD.encode(Sha256::digest(checkpoint_data))
}
fn digest_path(task_id: &str) -> std::path::PathBuf {
Path::new(BUCKET_META_PREFIX).join(format!("{task_id}_{RESUME_CHECKPOINT_DIGEST_FILE}"))
}
} }
-386
View File
@@ -1600,29 +1600,6 @@ async fn test_checkpoint_schema_v4_discarded_on_load() {
temp_dir.close().expect("remove schema test directory"); temp_dir.close().expect("remove schema test directory");
} }
#[tokio::test]
async fn unsigned_previous_checkpoint_schema_preserves_progress() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let mut legacy = ResumeCheckpoint::new(task_id.clone());
legacy.schema_version = CURRENT_CHECKPOINT_SCHEMA - 1;
legacy.update_position(2, 500);
legacy.add_processed_object("object".to_string());
legacy.integrity_digest = None;
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, serde_json::to_vec(&legacy).unwrap().into())
.await
.expect("write previous-schema checkpoint");
let manager = CheckpointManager::load_from_disk(disk, &task_id).await.unwrap();
let checkpoint = manager.get_checkpoint().await;
assert_eq!(checkpoint.schema_version, CURRENT_CHECKPOINT_SCHEMA);
assert_eq!(checkpoint.current_bucket_index, 2);
assert_eq!(checkpoint.current_object_index, 500);
assert!(checkpoint.processed_objects.contains("object"));
temp_dir.close().unwrap();
}
#[tokio::test] #[tokio::test]
async fn current_normal_resume_schema_preserves_progress() { async fn current_normal_resume_schema_preserves_progress() {
let (temp_dir, disk) = schema_test_disk().await; let (temp_dir, disk) = schema_test_disk().await;
@@ -1698,369 +1675,6 @@ async fn future_resume_and_checkpoint_schemas_are_rejected() {
temp_dir.close().expect("remove schema test directory"); temp_dir.close().expect("remove schema test directory");
} }
#[tokio::test]
async fn checkpoint_save_does_not_replace_a_non_empty_truncated_snapshot() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("create checkpoint manager");
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
let truncated = b"{\"schema_version\":5,\"task_id\":";
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, truncated.as_slice().into())
.await
.expect("write truncated checkpoint fixture");
let error = manager
.update_position(2, 7)
.await
.expect_err("a truncated checkpoint must fail closed during save");
assert!(error.to_string().contains("Existing checkpoint is corrupt"));
assert_eq!(
disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
.await
.expect("read truncated checkpoint fixture"),
truncated.as_slice()
);
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
temp_dir.close().expect("remove checkpoint save test directory");
}
#[tokio::test]
async fn checkpoint_save_does_not_replace_a_future_schema_snapshot() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("create checkpoint manager");
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
let mut future = ResumeCheckpoint::new(task_id.clone());
future.schema_version = CURRENT_CHECKPOINT_SCHEMA + 1;
let future_bytes = serde_json::to_vec(&future).expect("serialize future checkpoint fixture");
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, future_bytes.clone().into())
.await
.expect("write future checkpoint fixture");
let error = manager
.update_position(2, 7)
.await
.expect_err("a future schema must fail closed during save");
assert!(error.to_string().contains("Existing checkpoint schema"));
assert_eq!(
disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
.await
.expect("read future checkpoint fixture"),
future_bytes
);
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
temp_dir.close().expect("remove future schema test directory");
}
#[tokio::test]
async fn checkpoint_digest_rejects_same_length_progress_tampering() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("create checkpoint manager");
manager
.add_processed_object("victim-a".to_string())
.await
.expect("persist checkpoint progress");
manager.update_position(1, 1).await.expect("flush checkpoint progress");
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
let original = disk
.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
.await
.expect("read checkpoint fixture");
let tampered = original
.windows(b"victim-a".len())
.position(|window| window == b"victim-a")
.map(|index| {
let mut bytes = original.to_vec();
bytes[index..index + b"victim-a".len()].copy_from_slice(b"victim-b");
bytes
})
.expect("checkpoint should contain the processed object");
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, tampered.into())
.await
.expect("write tampered checkpoint fixture");
assert!(CheckpointManager::load_from_disk(disk.clone(), &task_id).await.is_err());
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
temp_dir.close().expect("remove digest test directory");
}
#[tokio::test]
async fn checkpoint_integrity_survives_missing_legacy_sidecar() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap();
manager.update_position(2, 9).await.unwrap();
let digest_path = format!("{BUCKET_META_PREFIX}/{task_id}_ahm_checkpoint.sha256");
delete_resume_file(&disk, Path::new(&digest_path)).await.unwrap();
let restored = CheckpointManager::load_from_disk(disk, &task_id).await.unwrap();
let checkpoint = restored.get_checkpoint().await;
assert_eq!(checkpoint.current_bucket_index, 2);
assert_eq!(checkpoint.current_object_index, 9);
assert!(checkpoint.integrity_digest.is_some());
temp_dir.close().unwrap();
}
#[tokio::test]
async fn checkpoint_integrity_survives_multi_object_reload() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap();
for index in 0..32 {
manager.add_processed_object(format!("processed-{index}")).await.unwrap();
manager.add_failed_object(format!("failed-{index}")).await.unwrap();
manager.add_skipped_object(format!("skipped-{index}")).await.unwrap();
}
manager.update_position(2, 9).await.unwrap();
CheckpointManager::load_from_disk(disk, &task_id)
.await
.expect("a healthy multi-object checkpoint must survive reload");
temp_dir.close().unwrap();
}
#[tokio::test]
async fn checkpoint_integrity_rejects_a_removed_embedded_digest() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone()).await.unwrap();
manager.update_position(2, 9).await.unwrap();
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
let bytes = disk
.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
.await
.expect("read checkpoint fixture");
let mut value: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
value["current_object_index"] = serde_json::json!(10);
value.as_object_mut().unwrap().remove("integrity_digest");
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, serde_json::to_vec(&value).unwrap().into())
.await
.expect("write tampered checkpoint fixture");
assert!(
CheckpointManager::load_from_disk(disk.clone(), &task_id).await.is_err(),
"a current checkpoint without its embedded digest must fail closed"
);
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
temp_dir.close().unwrap();
}
#[tokio::test]
async fn new_checkpoint_manager_rebuilds_an_empty_snapshot() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
disk.write_all(RUSTFS_META_BUCKET, &checkpoint_path, EcstoreDiskBytes::new())
.await
.expect("write empty checkpoint fixture");
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("a new manager must rebuild an empty checkpoint");
manager
.update_position(3, 11)
.await
.expect("rebuilt checkpoint must remain writable");
assert!(CheckpointManager::has_checkpoint(&disk, &task_id).await);
temp_dir.close().expect("remove empty checkpoint test directory");
}
#[tokio::test]
async fn deleted_checkpoint_is_not_recreated_by_an_old_manager() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("create checkpoint manager");
manager.cleanup().await.expect("delete checkpoint fixture");
let error = manager
.update_position(1, 2)
.await
.expect_err("an old manager must not resurrect a deleted checkpoint");
assert!(error.to_string().contains("removed after this manager saved it"));
assert!(!CheckpointManager::has_checkpoint(&disk, &task_id).await);
temp_dir.close().expect("remove deleted checkpoint test directory");
}
#[cfg(unix)]
#[tokio::test]
async fn checkpoint_cleanup_leaves_no_task_specific_lock_artifact() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("create checkpoint manager");
let lock_path = Path::new(BUCKET_META_PREFIX)
.join(format!("{task_id}_{RESUME_CHECKPOINT_FILE}"))
.with_extension("rustfs-cas.lock");
let lock_path = temp_dir.path().join(RUSTFS_META_BUCKET).join(lock_path);
manager.cleanup().await.expect("delete checkpoint fixture");
assert!(
!lock_path.exists(),
"successful checkpoint cleanup must not leave a task-specific lock artifact"
);
}
#[tokio::test]
async fn an_empty_blocked_marker_still_blocks_resume_selection() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let manager = CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("create checkpoint manager");
let blocked_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}");
disk.write_all(RUSTFS_META_BUCKET, &blocked_path, EcstoreDiskBytes::new())
.await
.expect("write empty blocked marker fixture");
assert!(CheckpointManager::is_blocked(&disk, &task_id).await);
assert!(CheckpointManager::is_resumable(&disk, &task_id).await.is_err());
// Recovery requires replacing/cleaning the snapshot, then removing the
// marker; ordinary selector retries are intentionally not an unlock path.
manager.cleanup().await.expect("clean blocked checkpoint");
assert!(!CheckpointManager::is_blocked(&disk, &task_id).await);
temp_dir.close().expect("remove empty blocked marker test directory");
}
#[tokio::test]
async fn resumable_selector_skips_healthy_tasks_with_blocked_markers() {
let (temp_dir, disk) = schema_test_disk().await;
let tasks = [
(ResumeUtils::generate_task_id(), EcstoreDiskBytes::new()),
(ResumeUtils::generate_task_id(), EcstoreDiskBytes::from_static(b"blocked")),
];
for (task_id, marker) in &tasks {
ResumeManager::new(
disk.clone(),
task_id.clone(),
"erasure_set".to_string(),
"pool_0_set_0".to_string(),
vec!["bucket".to_string()],
)
.await
.expect("create healthy resume state");
CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("create healthy checkpoint");
let checkpoint_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_FILE}");
let checkpoint_bytes = disk
.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
.await
.expect("read healthy checkpoint before blocking");
let marker_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}");
disk.write_all(RUSTFS_META_BUCKET, &marker_path, marker.clone())
.await
.expect("write blocked marker");
assert!(ResumeUtils::get_resumable_tasks(&disk).await.is_err());
assert_eq!(
disk.read_all(RUSTFS_META_BUCKET, &checkpoint_path)
.await
.expect("read healthy checkpoint after blocking"),
checkpoint_bytes
);
}
temp_dir.close().expect("remove blocked selector test directory");
}
#[tokio::test]
async fn stale_checkpoint_manager_cannot_overwrite_newer_progress() {
let (temp_dir, disk) = schema_test_disk().await;
let task_id = ResumeUtils::generate_task_id();
let first = CheckpointManager::new(disk.clone(), task_id.clone())
.await
.expect("create first checkpoint manager");
let second = CheckpointManager::load_from_disk(disk.clone(), &task_id)
.await
.expect("load second checkpoint manager");
second
.update_position(4, 20)
.await
.expect("persist newer checkpoint progress");
let error = first
.update_position(1, 3)
.await
.expect_err("stale checkpoint manager must not overwrite newer progress");
assert!(error.to_string().contains("newer progress"));
let persisted = CheckpointManager::load_from_disk(disk.clone(), &task_id)
.await
.expect("load newer checkpoint progress")
.get_checkpoint()
.await;
assert_eq!(persisted.current_bucket_index, 4);
assert_eq!(persisted.current_object_index, 20);
temp_dir.close().expect("remove stale manager test directory");
}
#[tokio::test]
async fn resumable_selector_isolates_future_and_corrupt_checkpoints() {
let (temp_dir, disk) = schema_test_disk().await;
let future_task = ResumeUtils::generate_task_id();
let corrupt_task = ResumeUtils::generate_task_id();
for task_id in [&future_task, &corrupt_task] {
ResumeManager::new(
disk.clone(),
task_id.to_string(),
"erasure_set".to_string(),
"pool_0_set_0".to_string(),
vec!["bucket".to_string()],
)
.await
.expect("create resumable state fixture");
}
let future_path = format!("{BUCKET_META_PREFIX}/{future_task}_{RESUME_CHECKPOINT_FILE}");
let mut future = ResumeCheckpoint::new(future_task.clone());
future.schema_version = CURRENT_CHECKPOINT_SCHEMA + 1;
let future_bytes = serde_json::to_vec(&future).expect("serialize future checkpoint fixture");
disk.write_all(RUSTFS_META_BUCKET, &future_path, future_bytes.clone().into())
.await
.expect("write future checkpoint fixture");
let corrupt_path = format!("{BUCKET_META_PREFIX}/{corrupt_task}_{RESUME_CHECKPOINT_FILE}");
let corrupt_bytes = b"{truncated";
disk.write_all(RUSTFS_META_BUCKET, &corrupt_path, corrupt_bytes.as_slice().into())
.await
.expect("write corrupt checkpoint fixture");
assert!(CheckpointManager::is_resumable(&disk, &future_task).await.is_err());
assert!(CheckpointManager::is_resumable(&disk, &corrupt_task).await.is_err());
assert!(ResumeUtils::get_resumable_tasks(&disk).await.is_err());
for (task_id, path, bytes) in [
(&future_task, future_path, future_bytes),
(&corrupt_task, corrupt_path, corrupt_bytes.to_vec()),
] {
assert_eq!(
disk.read_all(RUSTFS_META_BUCKET, &path)
.await
.expect("read isolated checkpoint bytes"),
bytes
);
let blocked_path = format!("{BUCKET_META_PREFIX}/{task_id}_{RESUME_CHECKPOINT_BLOCKED_FILE}");
assert!(
!disk
.read_all(RUSTFS_META_BUCKET, &blocked_path)
.await
.expect("read checkpoint blocked marker")
.is_empty()
);
}
temp_dir.close().expect("remove selector isolation test directory");
}
#[test] #[test]
fn test_persist_throttle_batches_until_threshold() { fn test_persist_throttle_batches_until_threshold() {
let mut throttle = PersistThrottle::new(); let mut throttle = PersistThrottle::new();
+1 -2
View File
@@ -21,7 +21,7 @@ use uuid::Uuid;
use super::super::{BUCKET_META_PREFIX, DiskError, DiskStore, HealDiskExt as _, RUSTFS_META_BUCKET}; use super::super::{BUCKET_META_PREFIX, DiskError, DiskStore, HealDiskExt as _, RUSTFS_META_BUCKET};
use super::replacement::{ReplacementPhase, ReplacementRecoveryRecord}; use super::replacement::{ReplacementPhase, ReplacementRecoveryRecord};
use super::{ use super::{
CheckpointManager, EVENT_HEAL_RESUME_STATE, LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, REPLACEMENT_COMPLETION_PROOF_FILE, EVENT_HEAL_RESUME_STATE, LOG_COMPONENT_HEAL, LOG_SUBSYSTEM_RESUME, REPLACEMENT_COMPLETION_PROOF_FILE,
REPLACEMENT_INTENT_FILE, RESUME_STATE_FILE, ResumeManager, ResumeStateFile, is_replacement_intent, path_to_str, REPLACEMENT_INTENT_FILE, RESUME_STATE_FILE, ResumeManager, ResumeStateFile, is_replacement_intent, path_to_str,
replacement_recovery_corruption_for_state_load, replacement_recovery_dir, validate_resume_task_id, replacement_recovery_corruption_for_state_load, replacement_recovery_dir, validate_resume_task_id,
}; };
@@ -67,7 +67,6 @@ impl ResumeUtils {
// Extract task ID from filename: {task_id}_ahm_resume_state.json // Extract task ID from filename: {task_id}_ahm_resume_state.json
if let Some(task_id) = entry.strip_suffix(&format!("_{RESUME_STATE_FILE}")) if let Some(task_id) = entry.strip_suffix(&format!("_{RESUME_STATE_FILE}"))
&& validate_resume_task_id(task_id).is_ok() && validate_resume_task_id(task_id).is_ok()
&& CheckpointManager::is_resumable(disk, task_id).await?
{ {
task_ids.push(task_id.to_string()); task_ids.push(task_id.to_string());
} }
+2 -7
View File
@@ -60,9 +60,7 @@ pub const REPLICATION_READ_ONLY_HISTORICAL_FIELDS: &[&str] = &[
"Destination.ReplicationTime", "Destination.ReplicationTime",
]; ];
// v2: disableProxy moved from unsupported to writable (per-target read-proxy pub const REMOTE_TARGET_CAPABILITY_CONTRACT_VERSION: u32 = 1;
// opt-out is accepted by set-remote-target and the `proxy` update op).
pub const REMOTE_TARGET_CAPABILITY_CONTRACT_VERSION: u32 = 2;
pub const REMOTE_TARGET_WRITABLE_FIELDS: &[&str] = &[ pub const REMOTE_TARGET_WRITABLE_FIELDS: &[&str] = &[
"sourcebucket", "sourcebucket",
@@ -85,12 +83,9 @@ pub const REMOTE_TARGET_WRITABLE_FIELDS: &[&str] = &[
// madmin default of 60s); the per-target health-check interval is not // madmin default of 60s); the per-target health-check interval is not
// yet applied — the heartbeat keeps its global env-configured interval. // yet applied — the heartbeat keeps its global env-configured interval.
"healthCheckDuration", "healthCheckDuration",
// Per-target read-proxy opt-out, consumed by the proxy-target selector
// (contract v2; previously only importable via MinIO bucket-targets.json).
"disableProxy",
]; ];
pub const REMOTE_TARGET_UNSUPPORTED_FIELDS: &[&str] = &["edge", "edgeSyncBeforeExpiry"]; pub const REMOTE_TARGET_UNSUPPORTED_FIELDS: &[&str] = &["disableProxy", "edge", "edgeSyncBeforeExpiry"];
#[derive(Debug, Clone, Serialize, Deserialize, Default)] #[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct ObjectOpts { pub struct ObjectOpts {
-33
View File
@@ -125,34 +125,6 @@ pub(crate) async fn read_config_with_revision<S: ScannerObjectIO>(
} }
} }
/// Read only the object revision without materializing its body.
pub(crate) async fn read_config_revision<S: ScannerObjectIO>(store: Arc<S>, path: &str) -> StorageResult<DataUsageCacheRevision> {
match store
.get_object_reader(
RUSTFS_META_BUCKET,
path,
None,
HeaderMap::new(),
&ObjectOptions {
no_lock: true,
..Default::default()
},
)
.await
{
Ok(reader) => reader
.object_info
.etag
.filter(|etag| !etag.is_empty())
.map(DataUsageCacheRevision::Etag)
.ok_or_else(|| StorageError::other(format!("scanner config object {path} has no ETag"))),
Err(Error::FileNotFound | Error::VolumeNotFound | Error::ObjectNotFound(_, _) | Error::BucketNotFound(_)) => {
Ok(DataUsageCacheRevision::Missing)
}
Err(err) => Err(err),
}
}
#[derive(Clone, Debug)] #[derive(Clone, Debug)]
pub(crate) struct DataUsageCacheRevisions { pub(crate) struct DataUsageCacheRevisions {
main: DataUsageCacheRevision, main: DataUsageCacheRevision,
@@ -174,11 +146,6 @@ pub static LEGACY_DATA_USAGE_OBJ_NAME_PATH: LazyLock<String> =
pub static DATA_USAGE_BLOOM_NAME_PATH: LazyLock<String> = pub static DATA_USAGE_BLOOM_NAME_PATH: LazyLock<String> =
LazyLock::new(|| format!("{BUCKET_META_PREFIX}{SLASH_SEPARATOR}{DATA_USAGE_BLOOM_NAME}")); LazyLock::new(|| format!("{BUCKET_META_PREFIX}{SLASH_SEPARATOR}{DATA_USAGE_BLOOM_NAME}"));
/// Durable companion object for a cycle-state object which cannot be decoded.
/// The primary object is deliberately never replaced or deleted by recovery.
pub static DATA_USAGE_BLOOM_RECOVERY_PATH: LazyLock<String> =
LazyLock::new(|| format!("{}.recovery-required.json", DATA_USAGE_BLOOM_NAME_PATH.as_str()));
pub static BACKGROUND_HEAL_INFO_PATH: LazyLock<String> = pub static BACKGROUND_HEAL_INFO_PATH: LazyLock<String> =
LazyLock::new(|| format!("{BUCKET_META_PREFIX}{SLASH_SEPARATOR}.background-heal.json")); LazyLock::new(|| format!("{BUCKET_META_PREFIX}{SLASH_SEPARATOR}.background-heal.json"));
@@ -74,7 +74,7 @@ impl DataUsageCache {
let loaded = Self::load_cache(store.clone(), name).await?; let loaded = Self::load_cache(store.clone(), name).await?;
let backup = match loaded.backup_revision { let backup = match loaded.backup_revision {
Some(revision) => Some(revision), Some(revision) => Some(revision),
None => match read_config_revision(store, &backup_path).await { None => match Self::revision_for_path(store, &backup_path).await {
Ok(revision) => Some(revision), Ok(revision) => Some(revision),
Err(err) => { Err(err) => {
counter!(METRIC_CACHE_BACKUP_REVISION_FAILURE_TOTAL).increment(1); counter!(METRIC_CACHE_BACKUP_REVISION_FAILURE_TOTAL).increment(1);
@@ -336,6 +336,33 @@ impl DataUsageCache {
} }
} }
async fn revision_for_path<S: ScannerObjectIO>(store: Arc<S>, path: &str) -> StorageResult<DataUsageCacheRevision> {
match store
.get_object_reader(
RUSTFS_META_BUCKET,
path,
None,
HeaderMap::new(),
&ObjectOptions {
no_lock: true,
..Default::default()
},
)
.await
{
Ok(reader) => reader
.object_info
.etag
.filter(|etag| !etag.is_empty())
.map(DataUsageCacheRevision::Etag)
.ok_or_else(|| StorageError::other(format!("scanner cache object {path} has no ETag"))),
Err(Error::FileNotFound | Error::VolumeNotFound | Error::ObjectNotFound(_, _) | Error::BucketNotFound(_)) => {
Ok(DataUsageCacheRevision::Missing)
}
Err(err) => Err(err),
}
}
pub(super) fn cache_save_timeout() -> Duration { pub(super) fn cache_save_timeout() -> Duration {
crate::runtime_config::scanner_cache_save_timeout() crate::runtime_config::scanner_cache_save_timeout()
} }
@@ -16,7 +16,7 @@ use super::persistence::DataUsageCacheLoadAttempt;
use super::*; use super::*;
use crate::storage_api::scanner_io::{HTTPRangeSpec, ObjectIO}; use crate::storage_api::scanner_io::{HTTPRangeSpec, ObjectIO};
use crate::{ScannerGetObjectReader, ScannerPutObjReader}; use crate::{ScannerGetObjectReader, ScannerPutObjReader};
use rustfs_data_usage::{ReplicationAllStats, ReplicationTargetUsage}; use rustfs_data_usage::{ReplicationAllStats, ReplicationStats};
use serde_json::Value; use serde_json::Value;
use std::io::Cursor; use std::io::Cursor;
use std::pin::Pin; use std::pin::Pin;
@@ -1636,7 +1636,7 @@ fn size_recursive_prunes_empty_and_preserves_threshold_replication_stats() {
replication_stats: Some(ReplicationAllStats { replication_stats: Some(ReplicationAllStats {
targets: HashMap::from([( targets: HashMap::from([(
"arn:test:threshold".to_string(), "arn:test:threshold".to_string(),
ReplicationTargetUsage { ReplicationStats {
after_threshold_count: 1, after_threshold_count: 1,
..Default::default() ..Default::default()
}, },
+1 -4
View File
@@ -75,10 +75,7 @@ pub use remote_scanner::{
}; };
pub use runtime_config::{apply_scanner_runtime_config, scanner_runtime_config_status, validate_scanner_runtime_config}; pub use runtime_config::{apply_scanner_runtime_config, scanner_runtime_config_status, validate_scanner_runtime_config};
pub use rustfs_common::last_minute; pub use rustfs_common::last_minute;
pub use scanner::{ pub use scanner::{ScannerCycleScheduleStatus, init_data_scanner, scanner_cycle_schedule_status, scanner_topology_digest};
ScannerCycleRecoveryMarker, ScannerCycleRecoveryStatus, ScannerCycleScheduleStatus, init_data_scanner,
reset_scanner_cycle_recovery, scanner_cycle_recovery_status, scanner_cycle_schedule_status, scanner_topology_digest,
};
pub use scanner_io::{ pub use scanner_io::{
ScannerDirtyUsageAckError, ScannerDirtyUsageState, acknowledge_dirty_usage_generation, clear_dirty_usage_bucket, ScannerDirtyUsageAckError, ScannerDirtyUsageState, acknowledge_dirty_usage_generation, clear_dirty_usage_bucket,
record_dirty_usage_bucket, record_scanner_maintenance_change, scanner_activity_epoch, scanner_dirty_usage_state, record_dirty_usage_bucket, record_scanner_maintenance_change, scanner_activity_epoch, scanner_dirty_usage_state,
+43 -84
View File
@@ -20,7 +20,7 @@ use std::sync::{Arc, LazyLock, RwLock};
use crate::data_usage_define::{ use crate::data_usage_define::{
BACKGROUND_HEAL_INFO_PATH, DATA_USAGE_BLOOM_NAME_PATH, DATA_USAGE_OBJ_NAME_PATH, DATA_USAGE_OBSERVED_OBJ_NAME_PATH, BACKGROUND_HEAL_INFO_PATH, DATA_USAGE_BLOOM_NAME_PATH, DATA_USAGE_OBJ_NAME_PATH, DATA_USAGE_OBSERVED_OBJ_NAME_PATH,
DataUsageCache, DataUsageCacheRevision, LEGACY_DATA_USAGE_OBJ_NAME_PATH, read_config_revision, read_config_with_revision, DataUsageCache, DataUsageCacheRevision, LEGACY_DATA_USAGE_OBJ_NAME_PATH, read_config_with_revision,
}; };
use crate::runtime_config::{ use crate::runtime_config::{
ScannerRuntimeConfig, ScannerRuntimeConfigSource, refresh_scanner_runtime_config_from_global, scanner_bitrot_cycle, ScannerRuntimeConfig, ScannerRuntimeConfigSource, refresh_scanner_runtime_config_from_global, scanner_bitrot_cycle,
@@ -54,7 +54,9 @@ use rustfs_config::{ENV_SCANNER_CYCLE, ENV_SCANNER_SPEED, ENV_SCANNER_START_DELA
use rustfs_data_usage::observed_data_usage_is_newer; use rustfs_data_usage::observed_data_usage_is_newer;
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use sha2::{Digest as _, Sha256}; use sha2::{Digest as _, Sha256};
use tokio::sync::{Notify, mpsc}; #[cfg(test)]
use tokio::sync::Notify;
use tokio::sync::mpsc;
use tokio::time::{Duration, Instant}; use tokio::time::{Duration, Instant};
use tokio_util::sync::CancellationToken; use tokio_util::sync::CancellationToken;
use tokio_util::task::AbortOnDropHandle; use tokio_util::task::AbortOnDropHandle;
@@ -102,13 +104,6 @@ const CLEAN_IDLE_BACKOFF_FACTOR: u32 = 2;
/// unavailable peer cannot drive a tight retry loop. /// unavailable peer cannot drive a tight retry loop.
const SCANNER_RETRY_BASE_INTERVAL: Duration = Duration::from_secs(5); const SCANNER_RETRY_BASE_INTERVAL: Duration = Duration::from_secs(5);
const SCANNER_RETRY_MAX_INTERVAL: Duration = Duration::from_secs(30 * 60); const SCANNER_RETRY_MAX_INTERVAL: Duration = Duration::from_secs(30 * 60);
/// A transient backend outage remains self-healing after the short retry
/// budget is exhausted, but the probe is intentionally sparse until storage
/// recovers or an operator reset wakes the scanner.
const SCANNER_CYCLE_RECOVERY_PAUSED_INTERVAL: Duration = Duration::from_secs(5 * 60);
/// Permanent recovery states still get a sparse status probe so a reset that
/// races the wait registration cannot leave the scanner asleep forever.
const SCANNER_CYCLE_RECOVERY_BLOCKED_PROBE_INTERVAL: Duration = Duration::from_secs(5 * 60);
const SCANNER_LEADER_LOCK_POLL_INTERVAL: Duration = Duration::from_secs(1); const SCANNER_LEADER_LOCK_POLL_INTERVAL: Duration = Duration::from_secs(1);
#[cfg(not(test))] #[cfg(not(test))]
const SCANNER_LOCK_LOSS_SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(30); const SCANNER_LOCK_LOSS_SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(30);
@@ -130,12 +125,6 @@ type ScannerCycleStatePersistTestHook = (u64, Arc<Notify>);
static SCANNER_CYCLE_STATE_PERSIST_TEST_HOOK: LazyLock<StdMutex<Option<ScannerCycleStatePersistTestHook>>> = static SCANNER_CYCLE_STATE_PERSIST_TEST_HOOK: LazyLock<StdMutex<Option<ScannerCycleStatePersistTestHook>>> =
LazyLock::new(|| StdMutex::new(None)); LazyLock::new(|| StdMutex::new(None));
static SCANNER_CYCLE_RECOVERY_WAKE: LazyLock<Notify> = LazyLock::new(Notify::new);
pub(super) fn notify_scanner_cycle_recovery_wake() {
SCANNER_CYCLE_RECOVERY_WAKE.notify_one();
}
#[cfg(test)] #[cfg(test)]
struct ScannerCycleStatePersistTestHookGuard; struct ScannerCycleStatePersistTestHookGuard;
@@ -587,21 +576,19 @@ pub async fn init_data_scanner(ctx: CancellationToken, storeapi: Arc<ECStore>) {
tokio::time::sleep(sleep_time).await; tokio::time::sleep(sleep_time).await;
} }
let mut transient_backoff = ScannerRetryBackoff::default();
let mut recovery_retry_count = 0_u32;
loop { loop {
if ctx_clone.is_cancelled() { if ctx_clone.is_cancelled() {
break; break;
} }
let run_result = run_data_scanner_with_maintenance_state( if let Err(e) = run_data_scanner_with_maintenance_state(
ctx_clone.clone(), ctx_clone.clone(),
storeapi_clone.clone(), storeapi_clone.clone(),
startup_features, startup_features,
startup_maintenance_generation, startup_maintenance_generation,
) )
.await; .await
if let Err(e) = &run_result { {
error!( error!(
target: "rustfs::scanner", target: "rustfs::scanner",
event = EVENT_SCANNER_CYCLE_STATE, event = EVENT_SCANNER_CYCLE_STATE,
@@ -612,52 +599,11 @@ pub async fn init_data_scanner(ctx: CancellationToken, storeapi: Arc<ECStore>) {
"Scanner runtime iteration failed" "Scanner runtime iteration failed"
); );
} }
let recovery_status = scanner_cycle_recovery_status();
if recovery_status.retryable {
recovery_retry_count = recovery_retry_count.saturating_add(1);
let _ = record_scanner_cycle_recovery_retry(recovery_retry_count);
} else {
recovery_retry_count = 0;
}
let recovery_status = scanner_cycle_recovery_status();
if recovery_status.state == "paused" {
transient_backoff.record_retryable_cycle(false);
tokio::select! {
_ = ctx_clone.cancelled() => break,
_ = SCANNER_CYCLE_RECOVERY_WAKE.notified() => {},
_ = tokio::time::sleep(SCANNER_CYCLE_RECOVERY_PAUSED_INTERVAL) => {},
}
recovery_retry_count = 0;
continue;
}
if !recovery_status.retryable
&& matches!(recovery_status.state.as_str(), "blocked" | "recovery-required" | "cleanup-pending")
{
transient_backoff.record_retryable_cycle(false);
tokio::select! {
_ = ctx_clone.cancelled() => break,
_ = SCANNER_CYCLE_RECOVERY_WAKE.notified() => {},
_ = tokio::time::sleep(SCANNER_CYCLE_RECOVERY_BLOCKED_PROBE_INTERVAL) => {},
}
continue;
}
let retry_delay = if recovery_status.retryable || run_result.is_err() {
transient_backoff.record_retryable_cycle(true);
transient_backoff
.retry_interval(scanner_cycle_interval())
.unwrap_or(SCANNER_RETRY_BASE_INTERVAL)
} else {
transient_backoff.record_retryable_cycle(false);
randomized_cycle_delay()
};
// Backoff before retrying after lock contention or scanner-level failures. // Backoff before retrying after lock contention or scanner-level failures.
// Keep this cancellation-aware so shutdown is not delayed by backoff sleep. // Keep this cancellation-aware so shutdown is not delayed by backoff sleep.
tokio::select! { tokio::select! {
_ = ctx_clone.cancelled() => break, _ = ctx_clone.cancelled() => break,
_ = SCANNER_CYCLE_RECOVERY_WAKE.notified() => {}, _ = tokio::time::sleep(randomized_cycle_delay()) => {}
_ = tokio::time::sleep(retry_delay) => {}
} }
} }
}); });
@@ -1660,22 +1606,40 @@ async fn run_data_scanner_with_maintenance_state(
observe_scanner_activity(&storeapi, distributed, &mut scanner_activity_seen).await; observe_scanner_activity(&storeapi, distributed, &mut scanner_activity_seen).await;
} }
let (mut cycle_info, mut leader_epoch, mut cycle_revision) = let (buf, mut cycle_revision) = match read_config_with_revision(storeapi.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()).await {
match load_scanner_cycle_state_for_startup(storeapi.clone()).await { Ok((buf, revision)) => (buf.unwrap_or_default(), revision),
ScannerCycleStateStartup::Ready { Err(err) => {
cycle, error!(
leader_epoch, target: "rustfs::scanner",
revision, event = EVENT_SCANNER_PERSIST_STATE,
} => (cycle, leader_epoch, revision), component = LOG_COMPONENT_SCANNER,
ScannerCycleStateStartup::Blocked => { subsystem = LOG_SUBSYSTEM_RUNTIME,
global_metrics().set_cycle(None).await; path = %&*DATA_USAGE_BLOOM_NAME_PATH,
return Ok(()); state = "revision_load_failed",
} error = %err,
ScannerCycleStateStartup::Transient(err) => { "Scanner cycle state revision load failed"
global_metrics().set_cycle(None).await; );
return Err(err); global_metrics().set_cycle(None).await;
} return Ok(());
}; }
};
let (mut cycle_info, mut leader_epoch) = match decode_scanner_cycle_state_for_startup(&buf) {
Ok(state) => state,
Err(err) => {
error!(
target: "rustfs::scanner",
event = EVENT_SCANNER_PERSIST_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_RUNTIME,
path = %&*DATA_USAGE_BLOOM_NAME_PATH,
state = "cycle_decode_failed",
error = %err,
"Scanner stopped because persisted cycle state is invalid"
);
global_metrics().set_cycle(None).await;
return Ok(());
}
};
let usage_floor = match persisted_usage_floor(storeapi.clone()).await { let usage_floor = match persisted_usage_floor(storeapi.clone()).await {
Ok(floor) => floor, Ok(floor) => floor,
Err(err) => { Err(err) => {
@@ -2255,12 +2219,7 @@ pub(crate) use activity::{
pub(crate) use activity::{ScannerCycleOutcome, scanner_cycle_outcome_with_pending_maintenance}; pub(crate) use activity::{ScannerCycleOutcome, scanner_cycle_outcome_with_pending_maintenance};
#[cfg(test)] #[cfg(test)]
pub(crate) use cycle_state::encode_scanner_cycle_fence_for_test; pub(crate) use cycle_state::encode_scanner_cycle_fence_for_test;
pub use cycle_state::{ pub(crate) use cycle_state::{current_scanner_leader_epoch, decode_persisted_scanner_cycle_fence};
ScannerCycleRecoveryMarker, ScannerCycleRecoveryStatus, reset_scanner_cycle_recovery, scanner_cycle_recovery_status,
};
pub(crate) use cycle_state::{
current_scanner_leader_epoch, decode_persisted_scanner_cycle_fence, load_scanner_cycle_state_for_startup,
};
pub use heal_info::{BackgroundHealInfo, read_background_heal_info, save_background_heal_info}; pub use heal_info::{BackgroundHealInfo, read_background_heal_info, save_background_heal_info};
pub use usage_store::store_data_usage_in_backend; pub use usage_store::store_data_usage_in_backend;
File diff suppressed because it is too large Load Diff
+1 -1
View File
@@ -196,7 +196,7 @@ pub(super) async fn claim_scanner_leadership(
if ctx.is_cancelled() { if ctx.is_cancelled() {
return false; return false;
} }
let Some(claimed_epoch) = persisted_epoch.checked_add(1).filter(|epoch| *epoch < u64::MAX) else { let Some(claimed_epoch) = persisted_epoch.checked_add(1) else {
error!( error!(
target: "rustfs::scanner", target: "rustfs::scanner",
event = EVENT_SCANNER_PERSIST_STATE, event = EVENT_SCANNER_PERSIST_STATE,
+5 -926
View File
@@ -15,12 +15,11 @@
use super::*; use super::*;
use crate::EcstoreResult; use crate::EcstoreResult;
use crate::{ use crate::{
DATA_USAGE_BLOOM_RECOVERY_PATH, Endpoint, EndpointServerPools, Endpoints, InstanceContext, PoolEndpoints, Endpoint, EndpointServerPools, Endpoints, InstanceContext, PoolEndpoints, ScannerGetObjectReader as GetObjectReader,
ScannerGetObjectReader as GetObjectReader, ScannerObjectInfo as ObjectInfo, ScannerObjectOptions as ObjectOptions, ScannerObjectInfo as ObjectInfo, ScannerObjectOptions as ObjectOptions, ScannerPutObjReader as PutObjReader,
ScannerPutObjReader as PutObjReader, init_bucket_metadata_sys_for_scanner_tests, init_ecstore_config_for_scanner_tests, init_bucket_metadata_sys_for_scanner_tests, init_ecstore_config_for_scanner_tests, init_local_disks_with_instance_ctx,
init_local_disks_with_instance_ctx,
}; };
use std::collections::{HashMap, HashSet}; use std::collections::HashMap;
use std::io::Cursor; use std::io::Cursor;
use std::task::Poll; use std::task::Poll;
use temp_env::{with_var, with_var_unset}; use temp_env::{with_var, with_var_unset};
@@ -118,15 +117,6 @@ async fn scanner_cycle_lock_fence_bounds_uncooperative_shutdown() {
assert!(cycle_ctx.is_cancelled()); assert!(cycle_ctx.is_cancelled());
} }
#[tokio::test]
async fn scanner_cycle_recovery_wake_survives_wait_registration_race() {
notify_scanner_cycle_recovery_wake();
tokio::time::timeout(Duration::from_secs(1), SCANNER_CYCLE_RECOVERY_WAKE.notified())
.await
.expect("recovery wake should retain a permit until the waiter registers");
}
struct ScannerDefaultSpeedGuard; struct ScannerDefaultSpeedGuard;
impl ScannerDefaultSpeedGuard { impl ScannerDefaultSpeedGuard {
@@ -161,7 +151,6 @@ impl Drop for ScannerDefaultCycleGuard {
struct MemoryConfigStore { struct MemoryConfigStore {
objects: Mutex<HashMap<String, Vec<u8>>>, objects: Mutex<HashMap<String, Vec<u8>>>,
revisions: Mutex<HashMap<String, u64>>, revisions: Mutex<HashMap<String, u64>>,
non_regular_objects: Mutex<HashSet<String>>,
fail_put_number: Mutex<HashMap<String, usize>>, fail_put_number: Mutex<HashMap<String, usize>>,
object_not_found_put_number: Mutex<HashMap<String, usize>>, object_not_found_put_number: Mutex<HashMap<String, usize>>,
error_after_commit_put_number: Mutex<HashMap<String, usize>>, error_after_commit_put_number: Mutex<HashMap<String, usize>>,
@@ -202,16 +191,12 @@ impl crate::storage_api::scanner_io::ObjectIO for MemoryConfigStore {
.get(&key) .get(&key)
.cloned() .cloned()
.ok_or(EcstoreError::FileNotFound)?; .ok_or(EcstoreError::FileNotFound)?;
let data_len = i64::try_from(data.len()).expect("memory test object length should fit in i64"); let revision = *self.revisions.lock().await.entry(key).or_insert(1);
let revision = *self.revisions.lock().await.entry(key.clone()).or_insert(1);
let is_dir = self.non_regular_objects.lock().await.contains(&key);
Ok(GetObjectReader { Ok(GetObjectReader {
stream: Box::new(Cursor::new(data)), stream: Box::new(Cursor::new(data)),
object_info: ObjectInfo { object_info: ObjectInfo {
etag: Some(format!("memory-{revision}")), etag: Some(format!("memory-{revision}")),
size: data_len,
is_dir,
..Default::default() ..Default::default()
}, },
buffered_body: None, buffered_body: None,
@@ -812,10 +797,6 @@ fn scanner_cycle_state_decodes_legacy_and_fenced_formats() {
let (fenced_cycle, fenced_epoch) = decode_scanner_cycle_state(&fenced).expect("fenced cycle state should decode"); let (fenced_cycle, fenced_epoch) = decode_scanner_cycle_state(&fenced).expect("fenced cycle state should decode");
assert_eq!(fenced_cycle.next, 13); assert_eq!(fenced_cycle.next, 13);
assert_eq!(fenced_epoch, 7); assert_eq!(fenced_epoch, 7);
let mut trailing = fenced;
trailing.push(0);
assert!(decode_scanner_cycle_state(&trailing).is_err());
} }
#[test] #[test]
@@ -842,840 +823,6 @@ fn scanner_startup_fails_closed_on_nonempty_corrupt_cycle_state() {
assert!(encode_scanner_cycle_state(&exhausted, 7).is_err()); assert!(encode_scanner_cycle_state(&exhausted, 7).is_err());
} }
#[tokio::test]
async fn corrupt_cycle_state_is_quarantined_once() {
let store = Arc::new(MemoryConfigStore::default());
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
store.objects.lock().await.insert(state_key.clone(), vec![1]);
store.revisions.lock().await.insert(state_key.clone(), 7);
assert!(matches!(
load_scanner_cycle_state_for_startup(store.clone()).await,
ScannerCycleStateStartup::Blocked
));
let marker_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str());
let marker_data = store
.objects
.lock()
.await
.get(&marker_key)
.cloned()
.expect("corrupt state must leave a durable recovery marker");
let marker: ScannerCycleRecoveryMarker = serde_json::from_slice(&marker_data).expect("marker should be valid JSON");
assert_eq!(marker.primary_revision, "memory-7");
assert_eq!(marker.path, DATA_USAGE_BLOOM_NAME_PATH.as_str());
assert_eq!(marker.quarantine_path, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str());
assert_eq!(marker.classification, "corrupt");
// A second startup sees the matching marker before consuming the poison body.
assert!(matches!(
load_scanner_cycle_state_for_startup(store.clone()).await,
ScannerCycleStateStartup::Blocked
));
// Replacing the primary object advances its revision; the stale marker must
// not quarantine the newer, valid state.
let cycle = CurrentCycle {
next: 9,
..Default::default()
};
let encoded = encode_scanner_cycle_state(&cycle, 3).expect("valid state should encode");
store.objects.lock().await.insert(state_key.clone(), encoded);
store.revisions.lock().await.insert(state_key, 8);
assert!(matches!(
load_scanner_cycle_state_for_startup(store).await,
ScannerCycleStateStartup::Ready {
cycle: CurrentCycle { next: 9, .. },
leader_epoch: 3,
..
}
));
}
#[tokio::test]
async fn empty_cycle_state_object_is_quarantined_as_corrupt() {
let store = Arc::new(MemoryConfigStore::default());
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
store.objects.lock().await.insert(state_key.clone(), Vec::new());
store.revisions.lock().await.insert(state_key, 6);
assert!(matches!(
load_scanner_cycle_state_for_startup(store).await,
ScannerCycleStateStartup::Blocked
));
assert_eq!(scanner_cycle_recovery_status().classification.as_deref(), Some("corrupt"));
assert!(
scanner_cycle_recovery_status()
.reason
.as_deref()
.is_some_and(|reason| reason.contains("empty"))
);
}
#[tokio::test]
async fn future_cycle_state_schema_is_recovery_required() {
let store = Arc::new(MemoryConfigStore::default());
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
let mut future = 17_u64.to_le_bytes().to_vec();
future.extend_from_slice(b"RSCYC999");
future.extend_from_slice(&4_u64.to_le_bytes());
future.extend_from_slice(&[0x90]);
store.objects.lock().await.insert(state_key.clone(), future);
store.revisions.lock().await.insert(state_key, 13);
assert!(matches!(
load_scanner_cycle_state_for_startup(store).await,
ScannerCycleStateStartup::Blocked
));
assert_eq!(scanner_cycle_recovery_status().classification.as_deref(), Some("future_schema"));
}
#[tokio::test]
async fn concurrent_leaders_cannot_quarantine_newer_cycle_state() {
let store = Arc::new(MemoryConfigStore::default());
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
store.objects.lock().await.insert(state_key.clone(), vec![1]);
store.revisions.lock().await.insert(state_key, 4);
let (first, second) = tokio::join!(
load_scanner_cycle_state_for_startup(store.clone()),
load_scanner_cycle_state_for_startup(store.clone()),
);
assert!(matches!(first, ScannerCycleStateStartup::Blocked));
assert!(matches!(second, ScannerCycleStateStartup::Blocked));
let marker_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str());
let marker_data = store
.objects
.lock()
.await
.get(&marker_key)
.cloned()
.expect("one contender must publish the recovery marker");
let marker: ScannerCycleRecoveryMarker = serde_json::from_slice(&marker_data).expect("marker should decode");
assert_eq!(marker.primary_revision, "memory-4");
}
#[tokio::test]
async fn cleanup_pending_marker_blocks_a_rewritten_primary_after_restart() {
let store = Arc::new(MemoryConfigStore::default());
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
let marker_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str());
let encoded = encode_scanner_cycle_state(
&CurrentCycle {
next: 12,
..Default::default()
},
8,
)
.expect("valid state should encode");
store.objects.lock().await.insert(state_key.clone(), encoded);
store.revisions.lock().await.insert(state_key, 22);
let marker = ScannerCycleRecoveryMarker {
schema_version: 1,
primary_revision: "memory-21".to_string(),
generation: 11,
leader_epoch: 7,
classification: "corrupt".to_string(),
first_detected_at_unix_secs: 1,
last_attempt_at_unix_secs: 2,
retry_count: 1,
reason: "reset in progress".to_string(),
path: DATA_USAGE_BLOOM_NAME_PATH.clone(),
quarantine_path: DATA_USAGE_BLOOM_RECOVERY_PATH.clone(),
state: "cleanup-pending".to_string(),
};
store
.objects
.lock()
.await
.insert(marker_key.clone(), serde_json::to_vec(&marker).expect("marker should encode"));
store.revisions.lock().await.insert(marker_key, 3);
assert!(matches!(
load_scanner_cycle_state_for_startup(store).await,
ScannerCycleStateStartup::Blocked
));
assert_eq!(scanner_cycle_recovery_status().state, "cleanup-pending");
}
#[test]
fn full_rescan_reset_accepts_unknown_marker_fields_without_trusting_cursor() {
let marker = br#"{
"schema_version": 99,
"primary_revision": "memory-7",
"generation": 9000,
"leader_epoch": 9000,
"classification": "new-future-classification",
"first_detected_at_unix_secs": 1,
"last_attempt_at_unix_secs": 2,
"retry_count": 9,
"reason": "future marker",
"path": "buckets/.bloomcycle.bin",
"quarantine_path": "buckets/.bloomcycle.bin.recovery-required.json",
"future_field": {"cursor": "untrusted"}
}"#;
let decoded =
super::cycle_state::decode_recovery_marker_for_reset(marker, &DataUsageCacheRevision::Etag("memory-3".to_string()))
.expect("full-rescan compatibility decoder should accept additive fields");
assert_eq!(decoded.primary_revision, "memory-7");
assert_eq!(decoded.classification, "future_schema");
assert_eq!(decoded.generation, 0);
assert_eq!(decoded.leader_epoch, 0);
assert_eq!(decoded.state, "blocked");
let malformed =
super::cycle_state::decode_recovery_marker_for_reset(b"{not-json", &DataUsageCacheRevision::Etag("memory-4".to_string()))
.expect("a full-rescan reset must recover even when the marker is malformed");
assert!(malformed.primary_revision.is_empty());
assert_eq!(malformed.classification, "future_schema");
}
#[tokio::test]
async fn full_rescan_reset_rebuilds_after_malformed_marker_without_trusting_cursor() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01])
.await
.expect("corrupt cycle state should be persisted");
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), br#"{not-json"#.to_vec())
.await
.expect("malformed marker should be persisted");
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.expect("full-rescan reset should recover malformed marker");
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("rebuilt cycle state should remain durable");
let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode");
assert_eq!(cycle.next, 0, "reset must use the verified usage floor, not marker cursor");
assert_eq!(leader_epoch, 1);
assert!(matches!(
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
Err(EcstoreError::ConfigNotFound)
));
}
#[tokio::test]
async fn full_rescan_reset_ignores_epoch_from_malformed_future_primary() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
let mut future_primary = vec![0; 24];
future_primary[8..16].copy_from_slice(b"RSCY9999");
future_primary[16..24].copy_from_slice(&u64::MAX.to_le_bytes());
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), future_primary)
.await
.expect("future cycle state should be persisted");
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), br#"{not-json"#.to_vec())
.await
.expect("malformed marker should be persisted");
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.expect("full-rescan reset should recover malformed future state");
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("rebuilt cycle state should remain durable");
let (_, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode");
assert_eq!(leader_epoch, 1, "invalid persisted bytes must not raise the recovery epoch");
assert!(matches!(
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
Err(EcstoreError::ConfigNotFound)
));
}
#[tokio::test]
async fn ecstore_exact_recovery_marker_delete_honors_etag() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"marker-v1".to_vec())
.await
.expect("initial recovery marker should be persisted");
let (_, stale_revision) = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
.await
.expect("initial marker revision should load");
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"marker-v2".to_vec())
.await
.expect("replacement recovery marker should be persisted");
let delete_result = store
.delete_config_object(
RUSTFS_META_BUCKET,
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
ObjectOptions {
http_preconditions: Some(stale_revision.preconditions()),
..Default::default()
},
)
.await;
assert!(matches!(delete_result, Err(EcstoreError::PreconditionFailed)));
assert_eq!(
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
.await
.expect("replacement marker should remain durable"),
b"marker-v2"
);
}
#[tokio::test]
async fn full_rescan_reset_rejects_corrupt_primary_under_stale_blocked_marker() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
let corrupt_primary = vec![0xff, 0x00, 0x01];
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), corrupt_primary.clone())
.await
.expect("corrupt cycle state should be persisted");
let (_, primary_revision) = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("primary revision should load");
let marker = ScannerCycleRecoveryMarker {
schema_version: 1,
primary_revision: "memory-stale".to_string(),
generation: 1,
leader_epoch: 1,
classification: "corrupt".to_string(),
first_detected_at_unix_secs: 1,
last_attempt_at_unix_secs: 2,
retry_count: 1,
reason: "blocked primary changed".to_string(),
path: DATA_USAGE_BLOOM_NAME_PATH.clone(),
quarantine_path: DATA_USAGE_BLOOM_RECOVERY_PATH.clone(),
state: "blocked".to_string(),
};
let marker_data = serde_json::to_vec(&marker).expect("blocked marker should encode");
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), marker_data.clone())
.await
.expect("blocked marker should be persisted");
assert!(
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.is_err(),
"a strict marker must fail closed when its primary revision changed"
);
assert_eq!(
read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("primary should remain readable"),
corrupt_primary
);
assert_eq!(
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
.await
.expect("blocked marker should remain durable"),
marker_data
);
assert!(!matches!(primary_revision, DataUsageCacheRevision::Missing));
}
#[tokio::test]
async fn full_rescan_reset_preserves_valid_primary_when_marker_is_malformed() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
let primary = CurrentCycle {
next: 42,
..Default::default()
};
let old_primary_data = encode_scanner_cycle_state(&primary, 7).expect("valid cycle state should encode");
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), old_primary_data.clone())
.await
.expect("valid cycle state should be persisted");
let (_, old_primary_revision) = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("primary state revision should load");
let old_usage = DataUsageInfo {
scanner_epoch: Some(7),
scanner_cycle: Some(41),
..Default::default()
};
let old_usage_data = serde_json::to_vec(&old_usage).expect("usage snapshot should encode");
save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), old_usage_data.clone())
.await
.expect("usage snapshot should be persisted");
let (_, old_usage_revision) = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("usage snapshot revision should load");
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec())
.await
.expect("malformed marker should be persisted");
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.expect("reset should clear a stale malformed marker");
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("valid primary should remain durable");
let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("primary cycle state should decode");
assert_eq!(cycle.next, 42, "reset must not regress an independently fenced primary");
assert_eq!(leader_epoch, 8, "reset must advance the preserved primary epoch");
let stale_primary_save = save_config_with_preconditions(
store.clone(),
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
old_primary_data,
old_primary_revision.preconditions(),
)
.await;
assert!(matches!(stale_primary_save, Err(EcstoreError::PreconditionFailed)));
let usage = read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("usage epoch fence should remain durable");
assert_eq!(
serde_json::from_slice::<DataUsageInfo>(&usage)
.expect("fenced usage should decode")
.scanner_epoch,
Some(8)
);
let stale_save = save_config_with_preconditions(
store.clone(),
DATA_USAGE_OBJ_NAME_PATH.as_str(),
old_usage_data,
old_usage_revision.preconditions(),
)
.await;
assert!(matches!(stale_save, Err(EcstoreError::PreconditionFailed)));
assert!(matches!(
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
Err(EcstoreError::ConfigNotFound)
));
}
#[tokio::test]
async fn full_rescan_reset_resumes_cleanup_pending_preserved_primary() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
let completed_at = Utc::now();
let primary = CurrentCycle {
current: 3,
next: 42,
cycle_completed: vec![completed_at],
started: completed_at,
};
save_config(
store.clone(),
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
encode_scanner_cycle_state(&primary, 7).expect("valid cycle state should encode"),
)
.await
.expect("valid cycle state should be persisted");
let usage = DataUsageInfo {
scanner_epoch: Some(7),
scanner_cycle: Some(41),
..Default::default()
};
save_config(
store.clone(),
DATA_USAGE_OBJ_NAME_PATH.as_str(),
serde_json::to_vec(&usage).expect("usage snapshot should encode"),
)
.await
.expect("usage snapshot should be persisted");
let marker = ScannerCycleRecoveryMarker {
schema_version: 1,
primary_revision: "memory-old".to_string(),
generation: 41,
leader_epoch: 7,
classification: "corrupt".to_string(),
first_detected_at_unix_secs: 1,
last_attempt_at_unix_secs: 2,
retry_count: 1,
reason: "reset in progress".to_string(),
path: DATA_USAGE_BLOOM_NAME_PATH.clone(),
quarantine_path: DATA_USAGE_BLOOM_RECOVERY_PATH.clone(),
state: "cleanup-pending".to_string(),
};
save_config(
store.clone(),
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
serde_json::to_vec(&marker).expect("marker should encode"),
)
.await
.expect("cleanup marker should be persisted");
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.expect("reset should resume a cleanup-pending preserved primary");
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("preserved cycle state should remain durable");
let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("cycle state should decode");
assert_eq!(cycle.current, 3, "cleanup retry must preserve the in-progress cursor");
assert_eq!(cycle.next, 42);
assert_eq!(cycle.cycle_completed, vec![completed_at]);
assert_eq!(cycle.started, completed_at);
assert_eq!(leader_epoch, 8);
let usage = read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("usage epoch fence should remain durable");
assert_eq!(
serde_json::from_slice::<DataUsageInfo>(&usage)
.expect("usage should decode")
.scanner_epoch,
Some(8)
);
assert!(matches!(
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
Err(EcstoreError::ConfigNotFound)
));
}
#[tokio::test]
async fn full_rescan_reset_rebuilds_oversized_regular_primary_with_malformed_marker() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0; 1024 * 1024 + 1])
.await
.expect("oversized cycle state should be persisted");
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec())
.await
.expect("malformed marker should be persisted");
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.expect("explicit full-rescan reset should replace an oversized regular primary");
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("rebuilt cycle state should remain durable");
let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode");
assert_eq!(cycle.next, 0);
assert_eq!(leader_epoch, 1);
assert!(matches!(
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
Err(EcstoreError::ConfigNotFound)
));
}
#[tokio::test]
async fn full_rescan_reset_rebuilds_oversized_primary_after_cleanup_marker() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0; 1024 * 1024 + 1])
.await
.expect("oversized cycle state should be persisted");
let (_, primary_revision) = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("primary revision should load");
let marker = ScannerCycleRecoveryMarker {
schema_version: 1,
primary_revision: match primary_revision {
DataUsageCacheRevision::Etag(etag) => etag,
DataUsageCacheRevision::Missing => panic!("primary revision should be present"),
},
generation: 1,
leader_epoch: 1,
classification: "corrupt".to_string(),
first_detected_at_unix_secs: 1,
last_attempt_at_unix_secs: 2,
retry_count: 1,
reason: "reset in progress".to_string(),
path: DATA_USAGE_BLOOM_NAME_PATH.clone(),
quarantine_path: DATA_USAGE_BLOOM_RECOVERY_PATH.clone(),
state: "cleanup-pending".to_string(),
};
save_config(
store.clone(),
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
serde_json::to_vec(&marker).expect("cleanup marker should encode"),
)
.await
.expect("cleanup marker should be persisted");
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.expect("cleanup retry should rebuild an oversized primary");
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("rebuilt cycle state should remain durable");
let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode");
assert_eq!(cycle.next, 0);
assert_eq!(leader_epoch, 1);
assert!(matches!(
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
Err(EcstoreError::ConfigNotFound)
));
}
#[tokio::test]
async fn full_rescan_reset_rebuilds_with_oversized_marker() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01])
.await
.expect("corrupt cycle state should be persisted");
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), vec![b'x'; 64 * 1024 + 1])
.await
.expect("oversized recovery marker should be persisted");
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.expect("full-rescan reset should recover an oversized marker");
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("rebuilt cycle state should remain durable");
let (_, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode");
assert_eq!(leader_epoch, 1);
assert!(matches!(
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
Err(EcstoreError::ConfigNotFound)
));
}
#[tokio::test]
async fn full_rescan_reset_rebuilds_with_empty_marker() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01])
.await
.expect("corrupt cycle state should be persisted");
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), Vec::new())
.await
.expect("empty recovery marker should be persisted");
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.expect("full-rescan reset should recover an empty marker");
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("rebuilt cycle state should remain durable");
let (_, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode");
assert_eq!(leader_epoch, 1);
assert!(matches!(
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
Err(EcstoreError::ConfigNotFound)
));
}
#[tokio::test]
async fn full_rescan_reset_keeps_cleanup_marker_when_preserved_epoch_is_exhausted() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
let primary = CurrentCycle {
next: 42,
..Default::default()
};
save_config(
store.clone(),
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
encode_scanner_cycle_state(&primary, u64::MAX).expect("valid cycle state should encode"),
)
.await
.expect("valid cycle state should be persisted");
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec())
.await
.expect("malformed marker should be persisted");
assert!(
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.is_err()
);
let marker = read_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
.await
.expect("cleanup marker should remain durable");
assert_eq!(
serde_json::from_slice::<ScannerCycleRecoveryMarker>(&marker)
.expect("cleanup marker should decode")
.state,
"cleanup-pending"
);
assert!(matches!(
load_scanner_cycle_state_for_startup(store).await,
ScannerCycleStateStartup::Blocked
));
}
#[tokio::test]
async fn full_rescan_reset_rejects_preserved_epoch_that_would_be_terminal() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
let primary = CurrentCycle {
next: 42,
..Default::default()
};
save_config(
store.clone(),
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
encode_scanner_cycle_state(&primary, u64::MAX - 1).expect("valid cycle state should encode"),
)
.await
.expect("valid cycle state should be persisted");
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec())
.await
.expect("malformed marker should be persisted");
assert!(
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.is_err(),
"reset must not persist the terminal leader epoch"
);
let marker = read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
.await
.expect("cleanup marker should remain durable");
assert_eq!(
serde_json::from_slice::<ScannerCycleRecoveryMarker>(&marker)
.expect("cleanup marker should decode")
.state,
"cleanup-pending"
);
}
#[tokio::test]
async fn full_rescan_reset_rejects_usage_floor_that_would_be_terminal() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), vec![0xff, 0x00, 0x01])
.await
.expect("corrupt cycle state should be persisted");
save_config(
store.clone(),
DATA_USAGE_OBJ_NAME_PATH.as_str(),
serde_json::to_vec(&DataUsageInfo {
scanner_epoch: Some(u64::MAX - 1),
..Default::default()
})
.expect("usage floor should encode"),
)
.await
.expect("usage floor should be persisted");
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec())
.await
.expect("malformed marker should be persisted");
assert!(
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.is_err(),
"reset must not persist the terminal leader epoch"
);
assert_eq!(
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
.await
.expect("recovery marker should remain durable"),
b"{not-json"
);
}
#[tokio::test]
async fn full_rescan_reset_rebuilds_empty_primary_with_malformed_marker() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), Vec::new())
.await
.expect("empty cycle state should be persisted");
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{not-json".to_vec())
.await
.expect("malformed marker should be persisted");
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.expect("explicit full-rescan reset should replace an empty primary");
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("rebuilt cycle state should remain durable");
let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode");
assert_eq!(cycle.next, 0);
assert_eq!(leader_epoch, 1);
assert!(matches!(
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
Err(EcstoreError::ConfigNotFound)
));
}
#[tokio::test]
async fn full_rescan_reset_rebuilds_when_primary_cycle_state_is_missing() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
let marker = ScannerCycleRecoveryMarker {
schema_version: 1,
primary_revision: "memory-missing".to_string(),
generation: u64::MAX,
leader_epoch: u64::MAX,
classification: "corrupt".to_string(),
first_detected_at_unix_secs: 1,
last_attempt_at_unix_secs: 2,
retry_count: 0,
reason: "missing primary".to_string(),
path: DATA_USAGE_BLOOM_NAME_PATH.clone(),
quarantine_path: DATA_USAGE_BLOOM_RECOVERY_PATH.clone(),
state: "blocked".to_string(),
};
save_config(
store.clone(),
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
serde_json::to_vec(&marker).expect("marker should encode"),
)
.await
.expect("marker should be persisted");
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.expect("full-rescan reset should recreate missing primary");
let state = read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("missing primary should be rebuilt");
let (cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("rebuilt cycle state should decode");
assert_eq!(cycle.next, 0);
assert_eq!(leader_epoch, 1);
assert!(matches!(
read_config(store, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()).await,
Err(EcstoreError::ConfigNotFound)
));
}
#[tokio::test]
async fn corrupt_cycle_state_rename_or_marker_failure_stays_recovery_required() {
let store = Arc::new(MemoryConfigStore::default());
let state_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
let marker_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str());
store.objects.lock().await.insert(state_key.clone(), vec![1]);
store.revisions.lock().await.insert(state_key, 9);
store.fail_put_number.lock().await.insert(marker_key, 1);
assert!(matches!(
load_scanner_cycle_state_for_startup(store.clone()).await,
ScannerCycleStateStartup::Transient(_)
));
let status = scanner_cycle_recovery_status();
assert_eq!(status.state, "recovery-required");
assert!(status.retryable);
assert!(
store
.objects
.lock()
.await
.contains_key(&memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str()))
);
}
#[tokio::test]
async fn oversized_or_symlinked_cycle_state_is_rejected() {
let store = Arc::new(MemoryConfigStore::default());
let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_NAME_PATH.as_str());
store.objects.lock().await.insert(key.clone(), vec![0; 1024 * 1024 + 1]);
store.revisions.lock().await.insert(key.clone(), 11);
assert!(matches!(
load_scanner_cycle_state_for_startup(store.clone()).await,
ScannerCycleStateStartup::Blocked
));
assert_eq!(scanner_cycle_recovery_status().classification.as_deref(), Some("corrupt"));
assert!(
scanner_cycle_recovery_status()
.reason
.as_deref()
.is_some_and(|reason| reason.contains("oversized"))
);
let marker_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str());
store.objects.lock().await.remove(&marker_key);
store.objects.lock().await.insert(key.clone(), vec![1]);
store.revisions.lock().await.insert(key.clone(), 12);
store.non_regular_objects.lock().await.insert(key);
// The object contract exposes a non-regular object as `is_dir`; local
// backends reject symlink/reparse entries before they become an object.
assert!(matches!(
load_scanner_cycle_state_for_startup(store).await,
ScannerCycleStateStartup::Blocked
));
}
#[tokio::test] #[tokio::test]
async fn scanner_startup_uses_primary_and_backup_usage_floor() { async fn scanner_startup_uses_primary_and_backup_usage_floor() {
let store = Arc::new(MemoryConfigStore::default()); let store = Arc::new(MemoryConfigStore::default());
@@ -1708,31 +855,6 @@ async fn scanner_startup_uses_primary_and_backup_usage_floor() {
assert_eq!(epoch, 11); assert_eq!(epoch, 11);
} }
#[tokio::test]
async fn scanner_usage_floor_ignores_older_backup_after_primary_epoch_fence() {
let store = Arc::new(MemoryConfigStore::default());
let backup_path = format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str());
for (path, epoch, cycle) in [(DATA_USAGE_OBJ_NAME_PATH.as_str(), 8, 100), (backup_path.as_str(), 7, 10_000)] {
store.objects.lock().await.insert(
memory_config_key(RUSTFS_META_BUCKET, path),
serde_json::to_vec(&DataUsageInfo {
scanner_epoch: Some(epoch),
scanner_cycle: Some(cycle),
..Default::default()
})
.expect("usage snapshot should encode"),
);
}
assert_eq!(
persisted_usage_floor(store).await.expect("usage floor should load"),
PersistedUsageFloor {
next_cycle: 101,
leader_epoch: 8,
}
);
}
#[test] #[test]
fn scanner_startup_treats_incomplete_usage_snapshot_as_cold() { fn scanner_startup_treats_incomplete_usage_snapshot_as_cold() {
let mut legacy = complete_usage_with_bucket_count(Some(std::time::SystemTime::now()), 1); let mut legacy = complete_usage_with_bucket_count(Some(std::time::SystemTime::now()), 1);
@@ -1865,15 +987,6 @@ async fn scanner_usage_floor_fails_closed_on_corrupt_or_exhausted_usage_state()
assert!(persisted_usage_floor(store.clone()).await.is_err()); assert!(persisted_usage_floor(store.clone()).await.is_err());
store.objects.lock().await.insert(
memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()),
br#"{}"#.to_vec(),
);
assert!(
persisted_usage_floor(store.clone()).await.is_err(),
"a structurally incomplete usage snapshot must not be treated as an empty floor"
);
store.objects.lock().await.insert( store.objects.lock().await.insert(
memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()), memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()),
serde_json::to_vec(&DataUsageInfo { serde_json::to_vec(&DataUsageInfo {
@@ -2131,22 +1244,6 @@ async fn test_leadership_claim_preserves_usage_epoch_floor_across_old_epoch_conf
assert_eq!(store.put_counts.lock().await.get(&key), Some(&3)); assert_eq!(store.put_counts.lock().await.get(&key), Some(&3));
} }
#[tokio::test]
async fn test_leadership_claim_rejects_terminal_epoch() {
let store = Arc::new(MemoryConfigStore::default());
let ctx = CancellationToken::new();
let mut revision = DataUsageCacheRevision::Missing;
let mut cycle = CurrentCycle {
next: 12,
..Default::default()
};
let mut persisted_epoch = u64::MAX - 1;
assert!(!claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch).await);
assert_eq!(persisted_epoch, u64::MAX - 1);
assert!(read_config(store, &DATA_USAGE_BLOOM_NAME_PATH).await.is_err());
}
#[tokio::test] #[tokio::test]
async fn test_leadership_claim_confirms_commit_after_returned_error() { async fn test_leadership_claim_confirms_commit_after_returned_error() {
let store = Arc::new(MemoryConfigStore::default()); let store = Arc::new(MemoryConfigStore::default());
@@ -3878,24 +2975,6 @@ fn superseded_retry_backoff_grows_from_the_default_cycle() {
} }
} }
#[tokio::test(start_paused = true)]
async fn corrupt_cycle_state_backoff_uses_virtual_clock() {
let mut backoff = ScannerRetryBackoff::default();
backoff.record_retryable_cycle(true);
let first_delay = backoff
.retry_interval(Duration::from_secs(60))
.expect("the first recovery retry should be scheduled");
assert_eq!(first_delay, Duration::from_secs(5));
let deadline = Instant::now() + first_delay;
assert!(Instant::now() < deadline);
tokio::time::advance(first_delay).await;
assert!(Instant::now() >= deadline);
backoff.record_retryable_cycle(true);
assert_eq!(backoff.retry_interval(Duration::from_secs(60)), Some(Duration::from_secs(10)));
}
#[test] #[test]
fn scanner_cycle_wait_plan_drives_growth_resets_and_bitrot_cap() { fn scanner_cycle_wait_plan_drives_growth_resets_and_bitrot_cap() {
let runtime_config = ScannerRuntimeConfig { let runtime_config = ScannerRuntimeConfig {
@@ -13,7 +13,7 @@
// limitations under the License. // limitations under the License.
use super::*; use super::*;
use rustfs_data_usage::{ReplicationAllStats, ReplicationTargetUsage}; use rustfs_data_usage::{ReplicationAllStats, ReplicationStats};
const TEST_PLAN_DIGEST: DataUsageScanPlanDigest = DataUsageScanPlanDigest([7; 32]); const TEST_PLAN_DIGEST: DataUsageScanPlanDigest = DataUsageScanPlanDigest([7; 32]);
@@ -271,7 +271,7 @@ fn completed_data_usage_info_flattens_nested_bucket_entries() {
replication_stats: Some(ReplicationAllStats { replication_stats: Some(ReplicationAllStats {
targets: HashMap::from([( targets: HashMap::from([(
"arn:target".to_string(), "arn:target".to_string(),
ReplicationTargetUsage { ReplicationStats {
replicated_size: 2048, replicated_size: 2048,
replicated_count: 2, replicated_count: 2,
..Default::default() ..Default::default()
-1
View File
@@ -76,7 +76,6 @@ 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,17 +218,6 @@ 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
@@ -352,19 +341,6 @@ 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,
-1
View File
@@ -126,7 +126,6 @@ mod tests {
let _list_remote_target_handler = replication::ListRemoteTargetHandler {}; let _list_remote_target_handler = replication::ListRemoteTargetHandler {};
let _remove_remote_target_handler = replication::RemoveRemoteTargetHandler {}; let _remove_remote_target_handler = replication::RemoveRemoteTargetHandler {};
let _scanner_status_handler = scanner::ScannerStatusHandler {}; let _scanner_status_handler = scanner::ScannerStatusHandler {};
let _scanner_cycle_state_reset_handler = scanner::ScannerCycleStateResetHandler {};
let _ilm_expiry_status_handler = scanner::IlmExpiryStatusHandler {}; let _ilm_expiry_status_handler = scanner::IlmExpiryStatusHandler {};
let _manual_transition_handler = ilm_transition::ManualTransitionRunHandler {}; let _manual_transition_handler = ilm_transition::ManualTransitionRunHandler {};
let _manual_transition_status_handler = ilm_transition::ManualTransitionJobStatusHandler {}; let _manual_transition_status_handler = ilm_transition::ManualTransitionJobStatusHandler {};
+7 -49
View File
@@ -73,8 +73,6 @@ enum TargetUpdateOp {
/// Connection group: credentials plus endpoint, target bucket, and TLS settings. /// Connection group: credentials plus endpoint, target bucket, and TLS settings.
Credentials, Credentials,
Sync, Sync,
/// Per-target read-proxy opt-out (`disableProxy`).
Proxy,
Bandwidth, Bandwidth,
Path, Path,
} }
@@ -83,13 +81,12 @@ fn parse_remote_target_update_ops(queries: &HashMap<String, String>) -> S3Result
const SUPPORTED_OPS: &[(&str, TargetUpdateOp)] = &[ const SUPPORTED_OPS: &[(&str, TargetUpdateOp)] = &[
("creds", TargetUpdateOp::Credentials), ("creds", TargetUpdateOp::Credentials),
("sync", TargetUpdateOp::Sync), ("sync", TargetUpdateOp::Sync),
("proxy", TargetUpdateOp::Proxy),
("bandwidth", TargetUpdateOp::Bandwidth), ("bandwidth", TargetUpdateOp::Bandwidth),
("path", TargetUpdateOp::Path), ("path", TargetUpdateOp::Path),
]; ];
// Present in the MinIO wire contract, but they drive target fields this // Present in the MinIO wire contract, but they drive target fields this
// version rejects as unsupported — fail loudly instead of silently ignoring. // version rejects as unsupported — fail loudly instead of silently ignoring.
const UNSUPPORTED_OPS: &[&str] = &["healthcheck", "edge", "edgeSyncBeforeExpiry"]; const UNSUPPORTED_OPS: &[&str] = &["proxy", "healthcheck", "edge", "edgeSyncBeforeExpiry"];
for key in UNSUPPORTED_OPS { for key in UNSUPPORTED_OPS {
if queries.get(*key).is_some_and(|value| value == "true") { if queries.get(*key).is_some_and(|value| value == "true") {
@@ -315,10 +312,11 @@ impl RemoteTargetRequest {
)); ));
} }
for (unsupported, configured) in REMOTE_TARGET_UNSUPPORTED_FIELDS for (unsupported, configured) in
.iter() REMOTE_TARGET_UNSUPPORTED_FIELDS
.copied() .iter()
.zip([self.edge, self.edge_sync_before_expiry]) .copied()
.zip([self.disable_proxy, self.edge, self.edge_sync_before_expiry])
{ {
if configured { if configured {
return Err(s3_error!( return Err(s3_error!(
@@ -704,7 +702,6 @@ impl Operation for SetRemoteTargetHandler {
target.deployment_id = remote_target.deployment_id.clone(); target.deployment_id = remote_target.deployment_id.clone();
} }
TargetUpdateOp::Sync => target.replication_sync = remote_target.replication_sync, TargetUpdateOp::Sync => target.replication_sync = remote_target.replication_sync,
TargetUpdateOp::Proxy => target.disable_proxy = remote_target.disable_proxy,
TargetUpdateOp::Bandwidth => target.bandwidth_limit = remote_target.bandwidth_limit, TargetUpdateOp::Bandwidth => target.bandwidth_limit = remote_target.bandwidth_limit,
TargetUpdateOp::Path => target.path = remote_target.path.clone(), TargetUpdateOp::Path => target.path = remote_target.path.clone(),
} }
@@ -1523,7 +1520,6 @@ mod tests {
("update", "true"), ("update", "true"),
("creds", "true"), ("creds", "true"),
("sync", "true"), ("sync", "true"),
("proxy", "true"),
("bandwidth", "true"), ("bandwidth", "true"),
("path", "true"), ("path", "true"),
])) ]))
@@ -1533,7 +1529,6 @@ mod tests {
vec![ vec![
TargetUpdateOp::Credentials, TargetUpdateOp::Credentials,
TargetUpdateOp::Sync, TargetUpdateOp::Sync,
TargetUpdateOp::Proxy,
TargetUpdateOp::Bandwidth, TargetUpdateOp::Bandwidth,
TargetUpdateOp::Path TargetUpdateOp::Path
] ]
@@ -2075,6 +2070,7 @@ mod tests {
("credentials.session_token", serde_json::json!("session-token")), ("credentials.session_token", serde_json::json!("session-token")),
("credentials.expiration", serde_json::json!("2026-01-01T00:00:00Z")), ("credentials.expiration", serde_json::json!("2026-01-01T00:00:00Z")),
("api", serde_json::json!("s3v2")), ("api", serde_json::json!("s3v2")),
("disableProxy", serde_json::json!(true)),
("edge", serde_json::json!(true)), ("edge", serde_json::json!(true)),
("edgeSyncBeforeExpiry", serde_json::json!(true)), ("edgeSyncBeforeExpiry", serde_json::json!(true)),
] { ] {
@@ -2304,44 +2300,6 @@ mod tests {
assert!(!REMOTE_TARGET_UNSUPPORTED_FIELDS.contains(&"healthCheckDuration")); assert!(!REMOTE_TARGET_UNSUPPORTED_FIELDS.contains(&"healthCheckDuration"));
} }
#[test]
fn remote_target_disable_proxy_is_declared_writable_edge_stays_unsupported() {
assert!(REMOTE_TARGET_WRITABLE_FIELDS.contains(&"disableProxy"));
assert!(!REMOTE_TARGET_UNSUPPORTED_FIELDS.contains(&"disableProxy"));
// edge sync has no implementation behind it — it must stay rejected.
assert!(REMOTE_TARGET_UNSUPPORTED_FIELDS.contains(&"edge"));
assert!(REMOTE_TARGET_UNSUPPORTED_FIELDS.contains(&"edgeSyncBeforeExpiry"));
}
#[test]
fn remote_target_create_accepts_disable_proxy() {
let mut request = valid_remote_target_request();
request["disableProxy"] = serde_json::json!(true);
let target = serde_json::from_value::<RemoteTargetRequest>(request)
.expect("request should deserialize")
.into_bucket_target()
.expect("disableProxy is a supported per-target read-proxy opt-out");
assert!(target.disable_proxy);
}
#[test]
fn update_body_with_proxy_op_toggles_disable_proxy_without_credentials() {
// Mirrors the other partial-update groups: a proxy-only update body may
// omit the connection fields entirely.
let body = serde_json::json!({
"arn": "arn:rustfs:replication:us-east-1:dep:target",
"type": "replication",
"disableProxy": true
});
let request: RemoteTargetRequest = serde_json::from_value(body).expect("partial update body should deserialize");
let target = request
.into_update_bucket_target(&[TargetUpdateOp::Proxy])
.expect("proxy-only update must not require credentials");
assert!(target.disable_proxy);
}
#[test] #[test]
fn remote_target_capability_fields_do_not_overlap() { fn remote_target_capability_fields_do_not_overlap() {
for field in REMOTE_TARGET_UNSUPPORTED_FIELDS { for field in REMOTE_TARGET_UNSUPPORTED_FIELDS {
+2 -95
View File
@@ -13,11 +13,8 @@
// limitations under the License. // limitations under the License.
use crate::admin::auth::authorize_admin_request; use crate::admin::auth::authorize_admin_request;
use crate::admin::handlers::supervise_admin_mutation;
use crate::admin::router::{AdminOperation, Operation, S3Router}; use crate::admin::router::{AdminOperation, Operation, S3Router};
use crate::admin::runtime_sources::{ use crate::admin::runtime_sources::current_scanner_metrics_report;
app_context_from_req, current_object_store_handle_for_context, current_scanner_metrics_report,
};
use crate::module_switches::{ENV_SCANNER_ENABLED, scanner_enabled_from_env}; use crate::module_switches::{ENV_SCANNER_ENABLED, scanner_enabled_from_env};
use crate::server::ADMIN_PREFIX; use crate::server::ADMIN_PREFIX;
use chrono::Utc; use chrono::Utc;
@@ -25,13 +22,11 @@ use http::{HeaderMap, HeaderValue};
use hyper::{Method, StatusCode}; use hyper::{Method, StatusCode};
use matchit::Params; use matchit::Params;
use rustfs_common::metrics::{ScannerLifecycleExpirySnapshot, ScannerMaintenanceControlSnapshot, ScannerMetricsReport}; use rustfs_common::metrics::{ScannerLifecycleExpirySnapshot, ScannerMaintenanceControlSnapshot, ScannerMetricsReport};
use rustfs_config::MAX_ADMIN_REQUEST_BODY_SIZE;
use rustfs_credentials::Credentials; use rustfs_credentials::Credentials;
use rustfs_policy::policy::action::{Action, AdminAction}; use rustfs_policy::policy::action::{Action, AdminAction};
use s3s::header::CONTENT_TYPE; use s3s::header::CONTENT_TYPE;
use s3s::{Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error}; use s3s::{Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error};
use serde::{Deserialize, Serialize}; use serde::Serialize;
use tokio_util::sync::CancellationToken;
const JSON_CONTENT_TYPE: &str = "application/json"; const JSON_CONTENT_TYPE: &str = "application/json";
@@ -43,13 +38,6 @@ struct ScannerStatusResponse {
metrics: ScannerMetricsReport, metrics: ScannerMetricsReport,
cycle_schedule: rustfs_scanner::ScannerCycleScheduleStatus, cycle_schedule: rustfs_scanner::ScannerCycleScheduleStatus,
runtime_config: rustfs_scanner::runtime_config::ScannerRuntimeConfigStatus, runtime_config: rustfs_scanner::runtime_config::ScannerRuntimeConfigStatus,
cycle_recovery: rustfs_scanner::ScannerCycleRecoveryStatus,
}
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
struct ScannerCycleResetRequest {
mode: String,
} }
#[derive(Debug, Serialize)] #[derive(Debug, Serialize)]
@@ -129,7 +117,6 @@ fn scanner_status_response(
metrics, metrics,
cycle_schedule, cycle_schedule,
runtime_config, runtime_config,
cycle_recovery: rustfs_scanner::scanner::scanner_cycle_recovery_status(),
} }
} }
@@ -157,11 +144,6 @@ pub fn register_scanner_route(r: &mut S3Router<AdminOperation>) -> std::io::Resu
format!("{ADMIN_PREFIX}/v3/scanner/status").as_str(), format!("{ADMIN_PREFIX}/v3/scanner/status").as_str(),
AdminOperation(&ScannerStatusHandler {}), AdminOperation(&ScannerStatusHandler {}),
)?; )?;
r.insert(
Method::POST,
format!("{ADMIN_PREFIX}/v3/scanner/cycle-state/reset").as_str(),
AdminOperation(&ScannerCycleStateResetHandler {}),
)?;
r.insert( r.insert(
Method::GET, Method::GET,
format!("{ADMIN_PREFIX}/v3/ilm/expiry/status").as_str(), format!("{ADMIN_PREFIX}/v3/ilm/expiry/status").as_str(),
@@ -181,13 +163,6 @@ async fn validate_scanner_status_request(req: &S3Request<Body>) -> S3Result<Cred
authorize_admin_request(req, vec![Action::AdminAction(AdminAction::ServerInfoAdminAction)]).await authorize_admin_request(req, vec![Action::AdminAction(AdminAction::ServerInfoAdminAction)]).await
} }
async fn validate_scanner_reset_request(req: &S3Request<Body>) -> S3Result<Credentials> {
if req.credentials.is_none() {
return Err(s3_error!(InvalidRequest, "missing credentials"));
}
authorize_admin_request(req, vec![Action::AdminAction(AdminAction::ConfigUpdateAdminAction)]).await
}
fn json_response(body: Vec<u8>) -> S3Result<S3Response<(StatusCode, Body)>> { fn json_response(body: Vec<u8>) -> S3Result<S3Response<(StatusCode, Body)>> {
let mut headers = HeaderMap::new(); let mut headers = HeaderMap::new();
let content_type = HeaderValue::from_str(JSON_CONTENT_TYPE) let content_type = HeaderValue::from_str(JSON_CONTENT_TYPE)
@@ -217,37 +192,6 @@ impl Operation for ScannerStatusHandler {
pub struct IlmExpiryStatusHandler {} pub struct IlmExpiryStatusHandler {}
pub struct ScannerCycleStateResetHandler {}
#[async_trait::async_trait]
impl Operation for ScannerCycleStateResetHandler {
async fn call(&self, mut req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let _cred = validate_scanner_reset_request(&req).await?;
let body = req
.input
.store_all_limited(MAX_ADMIN_REQUEST_BODY_SIZE)
.await
.map_err(|err| S3Error::with_message(S3ErrorCode::InvalidRequest, format!("invalid reset request body: {err}")))?;
let reset = serde_json::from_slice::<ScannerCycleResetRequest>(&body)
.map_err(|err| S3Error::with_message(S3ErrorCode::InvalidRequest, format!("invalid reset request body: {err}")))?;
if reset.mode != "full-rescan" {
return Err(S3Error::with_message(S3ErrorCode::InvalidRequest, "reset mode must be full-rescan"));
}
let context = app_context_from_req(&req)
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "storage layer not initialized"))?;
let store = current_object_store_handle_for_context(Some(context.as_ref()))
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "storage layer not initialized"))?;
supervise_admin_mutation("scanner cycle state reset", async move {
rustfs_scanner::scanner::reset_scanner_cycle_recovery(CancellationToken::new(), store)
.await
.map_err(|err| S3Error::with_message(S3ErrorCode::InternalError, err.to_string()))?;
Ok::<_, S3Error>(())
})
.await?;
json_response(br#"{"status":"reset","mode":"full-rescan"}"#.to_vec())
}
}
#[async_trait::async_trait] #[async_trait::async_trait]
impl Operation for IlmExpiryStatusHandler { impl Operation for IlmExpiryStatusHandler {
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> { async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
@@ -293,38 +237,6 @@ mod tests {
assert_eq!(err.message(), Some("missing credentials")); assert_eq!(err.message(), Some("missing credentials"));
} }
#[tokio::test]
async fn scanner_reset_gate_rejects_missing_credentials() {
let req = S3Request {
input: Body::from(String::new()),
method: Method::POST,
uri: http::Uri::from_static("/rustfs/admin/v3/scanner/cycle-state/reset"),
headers: HeaderMap::new(),
extensions: http::Extensions::new(),
credentials: None,
region: None,
service: None,
trailing_headers: None,
};
let err = validate_scanner_reset_request(&req)
.await
.expect_err("a reset request without credentials must be rejected");
assert_eq!(err.code(), &S3ErrorCode::InvalidRequest);
assert_eq!(err.message(), Some("missing credentials"));
}
#[test]
fn admin_reset_requires_full_rescan_or_verified_cursor() {
let full_rescan: ScannerCycleResetRequest =
serde_json::from_str(r#"{"mode":"full-rescan"}"#).expect("full rescan must be accepted");
assert_eq!(full_rescan.mode, "full-rescan");
let cursor: ScannerCycleResetRequest =
serde_json::from_str(r#"{"mode":"cursor"}"#).expect("mode validation belongs to the handler");
assert_ne!(cursor.mode, "full-rescan");
assert!(serde_json::from_str::<ScannerCycleResetRequest>(r#"{"mode":"full-rescan","cursor":"untrusted"}"#).is_err());
}
#[test] #[test]
fn scanner_disabled_reason_reports_startup_env_key() { fn scanner_disabled_reason_reports_startup_env_key() {
assert_eq!(scanner_disabled_reason(true), None); assert_eq!(scanner_disabled_reason(true), None);
@@ -392,11 +304,6 @@ mod tests {
assert_eq!(encoded["cycle_schedule"]["effective_interval_seconds"], 0); assert_eq!(encoded["cycle_schedule"]["effective_interval_seconds"], 0);
assert_eq!(encoded["cycle_schedule"]["clean_idle_backoff_enabled"], false); assert_eq!(encoded["cycle_schedule"]["clean_idle_backoff_enabled"], false);
assert_eq!(encoded["cycle_schedule"]["clean_idle_backoff_multiplier"], 1); assert_eq!(encoded["cycle_schedule"]["clean_idle_backoff_multiplier"], 1);
assert_eq!(encoded["cycle_recovery"]["state"], "healthy");
assert_eq!(
encoded["cycle_recovery"]["quarantine_path"],
rustfs_scanner::DATA_USAGE_BLOOM_RECOVERY_PATH.as_str()
);
} }
#[test] #[test]
+4 -21
View File
@@ -1262,9 +1262,7 @@ mod tests {
assert_eq!(response.summary.manual_transition_jobs.state, CapabilityState::Supported); assert_eq!(response.summary.manual_transition_jobs.state, CapabilityState::Supported);
assert_eq!(response.replication.contract_version, 1); assert_eq!(response.replication.contract_version, 1);
assert_eq!(response.replication.bucket_replication.contract_version, 1); assert_eq!(response.replication.bucket_replication.contract_version, 1);
// v2: disableProxy moved from unsupported to writable (per-target assert_eq!(response.replication.remote_targets.contract_version, 1);
// read-proxy opt-out reached the admin API).
assert_eq!(response.replication.remote_targets.contract_version, 2);
assert_eq!(response.replication.bucket_replication.status.state, CapabilityState::Supported); assert_eq!(response.replication.bucket_replication.status.state, CapabilityState::Supported);
assert_eq!(response.replication.remote_targets.status.state, CapabilityState::Supported); assert_eq!(response.replication.remote_targets.status.state, CapabilityState::Supported);
assert_eq!( assert_eq!(
@@ -1295,15 +1293,7 @@ mod tests {
.remote_targets .remote_targets
.fields .fields
.iter() .iter()
.any(|field| field.name == "disableProxy" && field.state == super::ReplicationFieldState::Supported) .any(|field| field.name == "disableProxy" && field.state == super::ReplicationFieldState::Unsupported)
);
assert!(
response
.replication
.remote_targets
.fields
.iter()
.any(|field| field.name == "edge" && field.state == super::ReplicationFieldState::Unsupported)
); );
assert!( assert!(
response response
@@ -1374,7 +1364,7 @@ mod tests {
assert_eq!(value["summary"]["manual_transition_jobs"]["state"], "supported"); assert_eq!(value["summary"]["manual_transition_jobs"]["state"], "supported");
assert_eq!(value["replication"]["contract_version"], 1); assert_eq!(value["replication"]["contract_version"], 1);
assert_eq!(value["replication"]["bucket_replication"]["contract_version"], 1); assert_eq!(value["replication"]["bucket_replication"]["contract_version"], 1);
assert_eq!(value["replication"]["remote_targets"]["contract_version"], 2); assert_eq!(value["replication"]["remote_targets"]["contract_version"], 1);
assert_eq!(value["replication"]["bucket_replication"]["status"]["state"], "supported"); assert_eq!(value["replication"]["bucket_replication"]["status"]["state"], "supported");
assert_eq!(value["replication"]["remote_targets"]["status"]["state"], "supported"); assert_eq!(value["replication"]["remote_targets"]["status"]["state"], "supported");
assert_eq!( assert_eq!(
@@ -1393,14 +1383,7 @@ mod tests {
.as_array() .as_array()
.expect("remote target fields should be an array") .expect("remote target fields should be an array")
.iter() .iter()
.any(|field| field["name"] == "disableProxy" && field["state"] == "supported") .any(|field| field["name"] == "disableProxy" && field["state"] == "unsupported")
);
assert!(
value["replication"]["remote_targets"]["fields"]
.as_array()
.expect("remote target fields should be an array")
.iter()
.any(|field| field["name"] == "edge" && field["state"] == "unsupported")
); );
assert!( assert!(
value["replication"]["remote_targets"]["fields"] value["replication"]["remote_targets"]["fields"]
-12
View File
@@ -428,12 +428,6 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[
admin(HttpMethod::Get, "/rustfs/admin/v3/config", CONFIG_UPDATE, RouteRiskLevel::High), admin(HttpMethod::Get, "/rustfs/admin/v3/config", CONFIG_UPDATE, RouteRiskLevel::High),
admin(HttpMethod::Put, "/rustfs/admin/v3/config", CONFIG_UPDATE, RouteRiskLevel::High), admin(HttpMethod::Put, "/rustfs/admin/v3/config", CONFIG_UPDATE, RouteRiskLevel::High),
admin(HttpMethod::Get, "/rustfs/admin/v3/scanner/status", SERVER_INFO, RouteRiskLevel::Sensitive), admin(HttpMethod::Get, "/rustfs/admin/v3/scanner/status", SERVER_INFO, RouteRiskLevel::Sensitive),
admin(
HttpMethod::Post,
"/rustfs/admin/v3/scanner/cycle-state/reset",
CONFIG_UPDATE,
RouteRiskLevel::High,
),
admin( admin(
HttpMethod::Get, HttpMethod::Get,
"/rustfs/admin/v3/ilm/expiry/status", "/rustfs/admin/v3/ilm/expiry/status",
@@ -2026,12 +2020,6 @@ mod tests {
assert_not_action(HttpMethod::Get, "/rustfs/admin/v3/ilm/expiry/status", SET_TIER); assert_not_action(HttpMethod::Get, "/rustfs/admin/v3/ilm/expiry/status", SET_TIER);
} }
#[test]
fn route_policy_requires_config_update_for_scanner_cycle_reset() {
assert_action(HttpMethod::Post, "/rustfs/admin/v3/scanner/cycle-state/reset", CONFIG_UPDATE);
assert_not_action(HttpMethod::Post, "/rustfs/admin/v3/scanner/cycle-state/reset", SERVER_INFO);
}
#[test] #[test]
fn route_policy_uses_tier_actions_for_transition_routes() { fn route_policy_uses_tier_actions_for_transition_routes() {
assert_action(HttpMethod::Post, "/rustfs/admin/v3/ilm/transition/run", SET_TIER); assert_action(HttpMethod::Post, "/rustfs/admin/v3/ilm/transition/run", SET_TIER);
@@ -243,7 +243,6 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
admin_route(Method::GET, "/v3/config"), admin_route(Method::GET, "/v3/config"),
admin_route(Method::PUT, "/v3/config"), admin_route(Method::PUT, "/v3/config"),
admin_route(Method::GET, "/v3/scanner/status"), admin_route(Method::GET, "/v3/scanner/status"),
admin_route(Method::POST, "/v3/scanner/cycle-state/reset"),
admin_route(Method::GET, "/v3/audit/target/list"), admin_route(Method::GET, "/v3/audit/target/list"),
admin_route_sample( admin_route_sample(
Method::PUT, Method::PUT,
@@ -880,7 +879,6 @@ fn test_register_routes_cover_representative_admin_paths() {
assert_route(&router, Method::GET, &admin_path("/v3/config")); assert_route(&router, Method::GET, &admin_path("/v3/config"));
assert_route(&router, Method::PUT, &admin_path("/v3/config")); assert_route(&router, Method::PUT, &admin_path("/v3/config"));
assert_route(&router, Method::GET, &admin_path("/v3/scanner/status")); assert_route(&router, Method::GET, &admin_path("/v3/scanner/status"));
assert_route(&router, Method::POST, &admin_path("/v3/scanner/cycle-state/reset"));
assert_route(&router, Method::GET, &admin_path("/v3/ilm/expiry/status")); assert_route(&router, Method::GET, &admin_path("/v3/ilm/expiry/status"));
assert_route(&router, Method::POST, &admin_path("/v3/ilm/transition/run")); assert_route(&router, Method::POST, &admin_path("/v3/ilm/transition/run"));
assert_route( assert_route(
@@ -1369,7 +1367,6 @@ fn test_admin_alias_paths_match_existing_admin_routes() {
(Method::GET, compat_admin_alias_path("/v3/config")), (Method::GET, compat_admin_alias_path("/v3/config")),
(Method::PUT, compat_admin_alias_path("/v3/config")), (Method::PUT, compat_admin_alias_path("/v3/config")),
(Method::GET, compat_admin_alias_path("/v3/scanner/status")), (Method::GET, compat_admin_alias_path("/v3/scanner/status")),
(Method::POST, compat_admin_alias_path("/v3/scanner/cycle-state/reset")),
(Method::GET, compat_admin_alias_path("/v3/ilm/expiry/status")), (Method::GET, compat_admin_alias_path("/v3/ilm/expiry/status")),
] { ] {
assert!( assert!(
+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_or_physical()), size: Some(object.get_actual_size().unwrap_or_default()),
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
+37 -234
View File
@@ -3969,55 +3969,6 @@ 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 {
@@ -8458,6 +8409,8 @@ 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
@@ -8475,23 +8428,32 @@ 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);
if !skip_stat { let (goi, source_missing) = 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(_) => {} Ok(res) => (res, false),
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)
@@ -8502,11 +8464,15 @@ 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
@@ -8523,8 +8489,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, accounting) = store let (dobjs, errs) = store
.delete_objects_with_tier_delete_journal_and_accounting(&bucket, object_to_delete.clone(), storage_delete_opts) .delete_objects_with_tier_delete_journal(&bucket, object_to_delete.clone(), storage_delete_opts)
.await; .await;
let _manager = get_concurrency_manager(); let _manager = get_concurrency_manager();
@@ -8549,16 +8515,17 @@ 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;
let committed_delete_marker = dobjs[i].delete_marker; if creates_delete_marker {
let delete_accounting = accounting.get(i).and_then(Option::as_ref); record_bucket_delete_marker_memory(&bucket).await;
let update = delete_memory_update( } else {
creates_delete_marker, let size = object_sizes[i].max(0) as u64;
committed_delete_marker, record_bucket_object_delete_memory(
delete_request_targets_current(object_to_delete[i].version_id), &bucket,
delete_accounting.and_then(|value| value.size), size,
delete_accounting.is_some_and(|value| value.removed_current_object), existing_object_infos[i].is_some() && object_to_delete[i].version_id.is_none(),
); )
apply_delete_memory_update(&bucket, update).await; .await;
}
} }
Err(error) => { Err(error) => {
delete_results[didx].error = Some(error); delete_results[didx].error = Some(error);
@@ -8836,24 +8803,12 @@ 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
// Prefix/force deletes and synthetic directory entries do not carry one if delete_creates_delete_marker(&opts) {
// committed object identity; leave their cache delta to reconciliation. record_bucket_delete_marker_memory(&bucket).await;
let update = if force_delete || obj_info.name.is_empty() || synthetic_version_id {
None
} else { } else {
// The storage commit returns this object's metadata while its record_bucket_object_delete_memory(&bucket, obj_info.size.max(0) as u64, opts.version_id.is_none()).await;
// 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 {
@@ -17906,158 +17861,6 @@ 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()
+1 -9
View File
@@ -72,11 +72,6 @@ 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,
@@ -1238,10 +1233,7 @@ 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, ecfs, object_utils, runtime}; pub(crate) use super::{bucket, data_usage, 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,
+1 -43
View File
@@ -296,10 +296,7 @@ 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),
// Compressed legacy objects may retain an unknown (-1) size: Some(v.get_actual_size().unwrap_or_default()),
// 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()
@@ -659,45 +656,6 @@ 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,8 +429,6 @@ 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,