Compare commits

..

1 Commits

Author SHA1 Message Date
overtrue e13049a846 feat(connect): emit durable heartbeats 2026-08-22 18:22:33 +08:00
44 changed files with 1952 additions and 3024 deletions
+9 -2
View File
@@ -45,6 +45,13 @@
# docker-capable self-hosted `dind-sm-standard-2` label was the alternative but
# has fewer cores and reintroduces fleet-state risk for no reliability gain.
# DISABLED. This workflow is switched off in the repository's Actions settings
# (state: disabled_manually) and does not run on any trigger, including its cron
# and workflow_dispatch. That state lives in GitHub's UI and is invisible when
# reading this file, which has already misled at least one audit — hence this
# banner. Re-enabling is a UI action; anyone doing so should first check that the
# workflow still matches the current CI layout. See rustfs/backlog#1603.
#
name: mint
on:
@@ -63,9 +70,9 @@ on:
- core
- full
mint-image:
description: "Mint image reference (empty = pinned default)"
description: "Mint image reference"
required: false
default: ""
default: "minio/mint:edge"
schedule:
# Weekly, after the Sunday s3-tests full sweep (starts 02:00 UTC, up to
# 3h) has finished, so the two never contend for the same runner pool.
+2 -2
View File
@@ -317,6 +317,8 @@ pub mod config {
}
pub mod data_usage {
#[cfg(feature = "test-util")]
pub use crate::data_usage::seed_bucket_usage_memory_for_test;
pub use crate::data_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,
@@ -328,8 +330,6 @@ pub mod data_usage {
remove_bucket_usage_from_backend, replace_bucket_usage_memory_from_info, store_compression_total_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 {
@@ -2855,7 +2855,7 @@ fn replicate_object_info_from_object_info(
.map(|v| OffsetDateTime::parse(&v, &Rfc3339).unwrap_or(OffsetDateTime::UNIX_EPOCH));
let mut rstate = oi.replication_state();
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 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();
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 {
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)),
version_id: object_info.version_id.map(|version_id| version_id.to_string()),
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,
content_type: object_info.content_type.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));
}
#[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]
fn replication_target_head_content_matches_compare_etag_only() {
let source = ObjectInfo {
+60 -63
View File
@@ -21,7 +21,7 @@ use crate::storage_api_contracts::{
bucket::{BucketInfo, BucketOperations, BucketOptions, DeleteBucketOptions, MakeBucketOptions},
list::{StorageListObjectVersionsInfo, StorageListObjectsV2Info, StorageObjectInfoOrErr, StorageWalkOptions},
multipart::{CompletePart, ListMultipartsInfo, ListPartsInfo, MultipartInfo, MultipartUploadResult, PartInfo},
object::{DeleteAccounting, DeletedObject, ObjectIO as _, ObjectOperations as _, ObjectToDelete},
object::{DeletedObject, ObjectIO as _, ObjectOperations as _, ObjectToDelete},
range::HTTPRangeSpec,
};
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]
impl crate::storage_api_contracts::object::ObjectIO for Sets {
type Error = Error;
@@ -715,8 +655,65 @@ impl crate::storage_api_contracts::object::ObjectOperations for Sets {
objects: Vec<ObjectToDelete>,
opts: ObjectOptions,
) -> (Vec<DeletedObject>, Vec<Option<Error>>) {
let (deleted, errors, _) = self.delete_objects_with_accounting(bucket, objects, opts).await;
(deleted, errors)
// Default return value
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))]
+8 -108
View File
@@ -1391,37 +1391,7 @@ impl BucketUsageAccumulator {
}
pub fn quota_object_size(object: &ObjectInfo) -> Result<u64, Error> {
// A compressed object may carry -1 while the transformed size is unknown
// (legacy streaming sentinel). In that case the persisted physical size
// is still a valid accounting floor; every other negative value is corrupt.
// An explicit negative `actual-size` metadata value is corrupt, however:
// the sentinel is only valid in the in-memory/object-part field written by
// the legacy streaming path, not as a persisted declared size.
let compressed = object.is_compressed();
if object.actual_size < -1 || (object.actual_size == -1 && !compressed) {
return Err(Error::PartMissingOrCorrupt);
}
if object
.parts
.iter()
.any(|part| part.actual_size < -1 || (part.actual_size < 0 && !compressed))
{
return Err(Error::PartMissingOrCorrupt);
}
let declared_actual_size = rustfs_utils::http::get_str(&object.user_defined, rustfs_utils::http::SUFFIX_ACTUAL_SIZE)
.filter(|value| !value.is_empty());
if declared_actual_size
.as_deref()
.and_then(|value| value.parse::<i64>().ok())
.is_some_and(|size| size < 0)
{
return Err(Error::PartMissingOrCorrupt);
}
let logical_size = match object.get_actual_size().map_err(Error::other)? {
size if size == -1 && compressed && declared_actual_size.is_none() => None,
size if size >= 0 => Some(u64::try_from(size).map_err(|_| Error::PartMissingOrCorrupt)?),
_ => return Err(Error::PartMissingOrCorrupt),
};
let logical_size = u64::try_from(object.get_actual_size().map_err(Error::other)?).map_err(|_| Error::PartMissingOrCorrupt)?;
let persisted_part_size = if object.parts.is_empty() {
u64::try_from(object.size).map_err(|_| Error::PartMissingOrCorrupt)?
} else {
@@ -1429,8 +1399,12 @@ pub fn quota_object_size(object: &ObjectInfo) -> Result<u64, Error> {
// Compressed streaming objects persist -1 when the transformed
// part size is unknown. The physical part size remains a valid
// quota floor; reject only non-negative values that overflow.
let actual_size = if part.actual_size == -1 {
0
let actual_size = if part.actual_size < 0 {
if object.is_compressed() {
0
} else {
return Err(Error::PartMissingOrCorrupt);
}
} else {
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)
})?
};
Ok(logical_size.unwrap_or(0).max(persisted_part_size))
Ok(logical_size.max(persisted_part_size))
}
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]
#[serial]
async fn live_bucket_usage_refreshes_are_coalesced_only_while_in_flight() {
+4 -34
View File
@@ -689,9 +689,6 @@ impl ObjectInfo {
}
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 {
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()))?;
return Ok(size);
}
if self.actual_size == -1 && self.parts.is_empty() {
return Ok(-1);
}
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);
}
let mut actual_size = 0;
self.parts.iter().for_each(|part| {
actual_size += part.actual_size;
});
if actual_size == 0 && 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)
}
/// 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 {
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,
},
namespace::NamespaceLocking as _,
object::{DeleteAccounting, DeletedObject, HTTPPreconditions, ObjectIO as _, ObjectOperations as _, ObjectToDelete},
object::{DeletedObject, HTTPPreconditions, ObjectIO as _, ObjectOperations as _, ObjectToDelete},
range::HTTPRangeSpec,
};
use crate::store::utils::is_reserved_or_invalid_bucket;
+4 -159
View File
@@ -45,7 +45,6 @@ use crate::bucket::replication::{
DeleteReplicationConfigSnapshot, ReplicationLifecycleBridge, ReplicationStatusType, VersionPurgeStatusType,
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::disk::{DataDirDeleteStatus, OldCurrentSize};
use crate::error::is_err_invalid_upload_id;
@@ -5656,18 +5655,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
objects: Vec<ObjectToDelete>,
opts: ObjectOptions,
) -> (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 accounting = vec![None; objects.len()];
let delete_config_snapshot = opts
.delete_replication_config_snapshot
.clone()
@@ -5757,7 +5745,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
*item = Some(Error::other(message.clone()));
}
}
return (del_objects, del_errs, accounting);
return (del_objects, del_errs);
}
},
}
@@ -5804,22 +5792,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
let source_missing = gerr
.as_ref()
.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
// client-facing identity, where `from_file_info` synthesizes
// `Some(Uuid::nil())` for a null version on a versioned or
@@ -5948,12 +5920,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
},
replication_state: vr.replication_state_internal.clone(),
..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
@@ -5999,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());
@@ -6237,16 +6204,7 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
}
}
// An accounting identity is actionable only when the delete result is
// successful. Never let a failed commit (including a partial quorum
// failure) reach the request-layer fast delta path.
for (index, err) in del_errs.iter().enumerate() {
if err.is_some() {
accounting[index] = None;
}
}
(del_objects, del_errs, accounting)
(del_objects, del_errs)
}
#[tracing::instrument(skip(self))]
@@ -6575,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);
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);
self.invalidate_get_object_metadata_cache(bucket, object).await;
Ok(obj_info)
@@ -7872,113 +7824,6 @@ mod replication_quota_safety_tests {
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]
async fn direct_put_cannot_persist_a_tiny_logical_size() {
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 crate::storage_api_contracts::range::HTTPRangeSpec;
pub(crate) use rustfs_storage_api::{
DeleteAccounting, DeletedObject, HTTPPreconditions, ObjectIO, ObjectLockDeleteOptions, ObjectLockRetentionOptions,
ObjectOperations, ObjectPreconditionError, ObjectPreconditionPart, ObjectPreconditionState, ObjectToDelete,
DeletedObject, HTTPPreconditions, ObjectIO, ObjectLockDeleteOptions, ObjectLockRetentionOptions, ObjectOperations,
ObjectPreconditionError, ObjectPreconditionPart, ObjectPreconditionState, ObjectToDelete,
};
pub(crate) trait EcstoreObjectIO:
+11 -54
View File
@@ -41,7 +41,7 @@ use crate::set_disk::{
};
use crate::storage_api_contracts::{
namespace::NamespaceLocking as _,
object::{DeleteAccounting, ObjectIO as _, ObjectOperations as _},
object::{ObjectIO as _, ObjectOperations as _},
};
use parking_lot::Mutex as ParkingMutex;
use rustfs_io_metrics::{
@@ -1216,14 +1216,6 @@ fn return_batch_delete_lock_error(objects: &[ObjectToDelete], err: Error) -> (Ve
(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> {
let mut object_names: Vec<&str> = objects.iter().map(|object| object.object_name.as_str()).collect();
object_names.sort_unstable();
@@ -2320,22 +2312,6 @@ impl ECStore {
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))]
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
@@ -2713,19 +2689,6 @@ impl ECStore {
opts: ObjectOptions,
tier_journal_api: Option<Arc<ECStore>>,
) -> (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
let objects: Vec<ObjectToDelete> = objects
.iter()
@@ -2738,7 +2701,6 @@ impl ECStore {
// Default return value
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());
for _ in 0..objects.len() {
@@ -2752,7 +2714,7 @@ impl ECStore {
} else {
match self.acquire_bucket_lifecycle_read_lock(bucket).await {
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() {
@@ -2764,21 +2726,21 @@ impl ECStore {
Err(err) => {
let message = err.to_string();
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)
&& 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) {
None
} else {
Some(match acquire_bucket_metadata_transaction_read_lock_in(&self.ctx, bucket).await {
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() {
@@ -2788,7 +2750,7 @@ impl ECStore {
let (state, incarnation_id, config_revision) =
match get_object_lock_config_and_incarnation_from_disk_in(&self.ctx, bucket).await {
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(
self.id,
@@ -2804,10 +2766,7 @@ impl ECStore {
if let (Some(expected), Some(current)) = (opts.expected_bucket_incarnation_id, current_bucket_incarnation_id)
&& expected != current
{
return return_batch_delete_lock_error_with_accounting(
objects.as_slice(),
StorageError::BucketNotFound(bucket.to_string()),
);
return return_batch_delete_lock_error(objects.as_slice(), StorageError::BucketNotFound(bucket.to_string()));
}
#[cfg(test)]
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 {
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());
@@ -2824,24 +2783,22 @@ impl ECStore {
if self.is_pool_rebalancing(pool.pool_idx).await {
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;
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 {
del_errs[idx] = None;
del_objects[idx] = dels[idx].clone();
accounting[idx] = pool_accounting[idx].clone();
break;
}
if del_errs[idx].is_none() {
del_errs[idx] = errs[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);
});
(del_objects, del_errs, accounting)
(del_objects, del_errs)
// let mut futures = Vec::with_capacity(objects.len());
-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)]
pub(crate) struct DataUsageCacheRevisions {
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> =
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> =
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 backup = match loaded.backup_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),
Err(err) => {
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 {
crate::runtime_config::scanner_cache_save_timeout()
}
+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 rustfs_common::last_minute;
pub use scanner::{
ScannerCycleRecoveryMarker, ScannerCycleRecoveryStatus, ScannerCycleScheduleStatus, init_data_scanner,
reset_scanner_cycle_recovery, scanner_cycle_recovery_status, scanner_cycle_schedule_status, scanner_topology_digest,
};
pub use scanner::{ScannerCycleScheduleStatus, init_data_scanner, scanner_cycle_schedule_status, scanner_topology_digest};
pub use scanner_io::{
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,
+43 -84
View File
@@ -20,7 +20,7 @@ use std::sync::{Arc, LazyLock, RwLock};
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,
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::{
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 serde::{Deserialize, Serialize};
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_util::sync::CancellationToken;
use tokio_util::task::AbortOnDropHandle;
@@ -102,13 +104,6 @@ const CLEAN_IDLE_BACKOFF_FACTOR: u32 = 2;
/// unavailable peer cannot drive a tight retry loop.
const SCANNER_RETRY_BASE_INTERVAL: Duration = Duration::from_secs(5);
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);
#[cfg(not(test))]
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>>> =
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)]
struct ScannerCycleStatePersistTestHookGuard;
@@ -587,21 +576,19 @@ pub async fn init_data_scanner(ctx: CancellationToken, storeapi: Arc<ECStore>) {
tokio::time::sleep(sleep_time).await;
}
let mut transient_backoff = ScannerRetryBackoff::default();
let mut recovery_retry_count = 0_u32;
loop {
if ctx_clone.is_cancelled() {
break;
}
let run_result = run_data_scanner_with_maintenance_state(
if let Err(e) = run_data_scanner_with_maintenance_state(
ctx_clone.clone(),
storeapi_clone.clone(),
startup_features,
startup_maintenance_generation,
)
.await;
if let Err(e) = &run_result {
.await
{
error!(
target: "rustfs::scanner",
event = EVENT_SCANNER_CYCLE_STATE,
@@ -612,52 +599,11 @@ pub async fn init_data_scanner(ctx: CancellationToken, storeapi: Arc<ECStore>) {
"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.
// Keep this cancellation-aware so shutdown is not delayed by backoff sleep.
tokio::select! {
_ = ctx_clone.cancelled() => break,
_ = SCANNER_CYCLE_RECOVERY_WAKE.notified() => {},
_ = tokio::time::sleep(retry_delay) => {}
_ = tokio::time::sleep(randomized_cycle_delay()) => {}
}
}
});
@@ -1660,22 +1606,40 @@ async fn run_data_scanner_with_maintenance_state(
observe_scanner_activity(&storeapi, distributed, &mut scanner_activity_seen).await;
}
let (mut cycle_info, mut leader_epoch, mut cycle_revision) =
match load_scanner_cycle_state_for_startup(storeapi.clone()).await {
ScannerCycleStateStartup::Ready {
cycle,
leader_epoch,
revision,
} => (cycle, leader_epoch, revision),
ScannerCycleStateStartup::Blocked => {
global_metrics().set_cycle(None).await;
return Ok(());
}
ScannerCycleStateStartup::Transient(err) => {
global_metrics().set_cycle(None).await;
return Err(err);
}
};
let (buf, mut cycle_revision) = match read_config_with_revision(storeapi.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()).await {
Ok((buf, revision)) => (buf.unwrap_or_default(), revision),
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 = "revision_load_failed",
error = %err,
"Scanner cycle state revision load failed"
);
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 {
Ok(floor) => floor,
Err(err) => {
@@ -2255,12 +2219,7 @@ pub(crate) use activity::{
pub(crate) use activity::{ScannerCycleOutcome, scanner_cycle_outcome_with_pending_maintenance};
#[cfg(test)]
pub(crate) use cycle_state::encode_scanner_cycle_fence_for_test;
pub use cycle_state::{
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(crate) use cycle_state::{current_scanner_leader_epoch, decode_persisted_scanner_cycle_fence};
pub use heal_info::{BackgroundHealInfo, read_background_heal_info, save_background_heal_info};
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() {
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!(
target: "rustfs::scanner",
event = EVENT_SCANNER_PERSIST_STATE,
+5 -926
View File
@@ -15,12 +15,11 @@
use super::*;
use crate::EcstoreResult;
use crate::{
DATA_USAGE_BLOOM_RECOVERY_PATH, Endpoint, EndpointServerPools, Endpoints, InstanceContext, PoolEndpoints,
ScannerGetObjectReader as GetObjectReader, ScannerObjectInfo as ObjectInfo, ScannerObjectOptions as ObjectOptions,
ScannerPutObjReader as PutObjReader, init_bucket_metadata_sys_for_scanner_tests, init_ecstore_config_for_scanner_tests,
init_local_disks_with_instance_ctx,
Endpoint, EndpointServerPools, Endpoints, InstanceContext, PoolEndpoints, ScannerGetObjectReader as GetObjectReader,
ScannerObjectInfo as ObjectInfo, ScannerObjectOptions as ObjectOptions, ScannerPutObjReader as PutObjReader,
init_bucket_metadata_sys_for_scanner_tests, init_ecstore_config_for_scanner_tests, init_local_disks_with_instance_ctx,
};
use std::collections::{HashMap, HashSet};
use std::collections::HashMap;
use std::io::Cursor;
use std::task::Poll;
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());
}
#[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;
impl ScannerDefaultSpeedGuard {
@@ -161,7 +151,6 @@ impl Drop for ScannerDefaultCycleGuard {
struct MemoryConfigStore {
objects: Mutex<HashMap<String, Vec<u8>>>,
revisions: Mutex<HashMap<String, u64>>,
non_regular_objects: Mutex<HashSet<String>>,
fail_put_number: Mutex<HashMap<String, usize>>,
object_not_found_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)
.cloned()
.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.clone()).or_insert(1);
let is_dir = self.non_regular_objects.lock().await.contains(&key);
let revision = *self.revisions.lock().await.entry(key).or_insert(1);
Ok(GetObjectReader {
stream: Box::new(Cursor::new(data)),
object_info: ObjectInfo {
etag: Some(format!("memory-{revision}")),
size: data_len,
is_dir,
..Default::default()
},
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");
assert_eq!(fenced_cycle.next, 13);
assert_eq!(fenced_epoch, 7);
let mut trailing = fenced;
trailing.push(0);
assert!(decode_scanner_cycle_state(&trailing).is_err());
}
#[test]
@@ -842,840 +823,6 @@ fn scanner_startup_fails_closed_on_nonempty_corrupt_cycle_state() {
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]
async fn scanner_startup_uses_primary_and_backup_usage_floor() {
let store = Arc::new(MemoryConfigStore::default());
@@ -1708,31 +855,6 @@ async fn scanner_startup_uses_primary_and_backup_usage_floor() {
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]
fn scanner_startup_treats_incomplete_usage_snapshot_as_cold() {
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());
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(
memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()),
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));
}
#[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]
async fn test_leadership_claim_confirms_commit_after_returned_error() {
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]
fn scanner_cycle_wait_plan_drives_growth_resets_and_bitrot_cap() {
let runtime_config = ScannerRuntimeConfig {
-1
View File
@@ -76,7 +76,6 @@ pub use bucket::{BucketInfo, BucketOperations, BucketOptions, DeleteBucketOption
pub use capability::{CapabilitySnapshotError, CapabilityState, CapabilityStatus};
pub use error::{StorageErrorCode, StorageResult};
pub use multipart::{CompletePart, ListMultipartsInfo, ListPartsInfo, MultipartInfo, MultipartUploadResult, PartInfo};
pub use object::DeleteAccounting;
pub use object::ObjectLockDeleteOptions;
pub use object::{DeletedObject, ObjectToDelete};
pub use object::{ExpirationOptions, TransitionedObject};
-24
View File
@@ -218,17 +218,6 @@ pub struct DeletedObject {
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 {
pub fn version_purge_status(&self) -> VersionPurgeStatusType {
self.replication_state
@@ -352,19 +341,6 @@ pub trait ObjectOperations: Send + Sync + fmt::Debug {
objects: Vec<Self::ObjectToDelete>,
opts: Self::ObjectOptions,
) -> (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(
&self,
bucket: &str,
+1 -1
View File
@@ -25,7 +25,7 @@
},
{
"name": "heartbeat",
"status": "reserved",
"status": "populated",
"purpose": "Heartbeat payloads, Connect receive time, and freshness window behavior."
},
{
@@ -0,0 +1,5 @@
975c1ca53eefeef6766a6fc0b3d3281f7408255342b0686e5e2aee5ad055414c duplicate.json
963529a38a02849c6c2acc6d72668dca9f63218b49c89fae41a451b584850411 overflow.json
e3adeee1c8a19aa17e70894896fb79c072e3785bea3611b93c11e79f039ed5af stale.json
35b9cebd8525389a701e8fe69fbe96407bcb31aa28392fe95babf4a4886985ad unknown.json
37941735dbd6ad3d238258a7b2cae6f0b3aa0ecaae1d8817817c3d718d11d633 valid.json
@@ -0,0 +1,9 @@
{
"protocolVersion": "v1",
"fixtureSet": "heartbeat",
"fixture": "duplicate",
"description": "An exact requestId replay returns the first result and creates no second heartbeat.",
"first": {"requestId": "550e8400-e29b-41d4-a716-446655440000", "sequence": 42},
"replay": {"requestId": "550e8400-e29b-41d4-a716-446655440000", "sequence": 42},
"expected": {"decision": "DUPLICATE", "heartbeatWrites": 1, "events": 1, "sameResponse": true}
}
@@ -0,0 +1,11 @@
{
"protocolVersion": "v1",
"fixtureSet": "heartbeat",
"fixture": "overflow",
"description": "Values beyond frozen bounds are rejected before persistence.",
"vectors": [
{"field": "sequence", "value": 9007199254740992, "maximum": 9007199254740991},
{"field": "coarseNodeSummary.total", "value": 4097, "maximum": 4096}
],
"expected": {"decision": "REJECT", "httpStatus": 422, "status": "INVALID_ARGUMENT"}
}
@@ -0,0 +1,9 @@
{
"protocolVersion": "v1",
"fixtureSet": "heartbeat",
"fixture": "stale",
"description": "A lower heartbeat sequence is retained as history and cannot replace the current projection.",
"head": {"requestId": "550e8400-e29b-41d4-a716-446655440000", "sequence": 42},
"late": {"requestId": "7c4d2e10-9f83-4a5b-b6c7-d8e9f0a1b2c3", "sequence": 9},
"expected": {"decision": "ACCEPT_HISTORY", "currentSequence": 42, "historySequence": 9}
}
@@ -0,0 +1,18 @@
{
"protocolVersion": "v1",
"fixtureSet": "heartbeat",
"fixture": "unknown",
"description": "Unknown optional members and capabilities are accepted, discarded before hashing, and never stored or echoed.",
"requestAdditions": {
"telemetryProfile": "extended",
"authorization": "Bearer non-functional-example",
"capabilities": ["heartbeat", "future.capability"],
"coarseNodeSummary": {"rackNames": ["customer-rack"]}
},
"expected": {
"decision": "ACCEPT",
"storedCapabilities": ["heartbeat"],
"discarded": ["authorization", "future.capability", "telemetryProfile", "coarseNodeSummary.rackNames"],
"echoed": []
}
}
@@ -0,0 +1,21 @@
{
"protocolVersion": "v1",
"fixtureSet": "heartbeat",
"fixture": "valid",
"description": "A bounded L0 heartbeat. clientTime is advisory; Connect's receivedAt is online authority.",
"request": {
"protocolVersion": "v1",
"requestId": "550e8400-e29b-41d4-a716-446655440000",
"agentVersion": "rustfs-agent/1.19.4",
"capabilities": ["heartbeat", "inventory"],
"sequence": 42,
"clientTime": "2026-08-22T01:02:03Z",
"coarseNodeSummary": {"total": 8, "healthy": 7, "degraded": 1}
},
"expected": {
"decision": "ACCEPT",
"acceptedVersion": "v1",
"responseFields": ["serverTime", "acceptedVersion", "capabilityHints"],
"onlineAuthority": "serverTime"
}
}
-1
View File
@@ -126,7 +126,6 @@ mod tests {
let _list_remote_target_handler = replication::ListRemoteTargetHandler {};
let _remove_remote_target_handler = replication::RemoveRemoteTargetHandler {};
let _scanner_status_handler = scanner::ScannerStatusHandler {};
let _scanner_cycle_state_reset_handler = scanner::ScannerCycleStateResetHandler {};
let _ilm_expiry_status_handler = scanner::IlmExpiryStatusHandler {};
let _manual_transition_handler = ilm_transition::ManualTransitionRunHandler {};
let _manual_transition_status_handler = ilm_transition::ManualTransitionJobStatusHandler {};
+2 -95
View File
@@ -13,11 +13,8 @@
// limitations under the License.
use crate::admin::auth::authorize_admin_request;
use crate::admin::handlers::supervise_admin_mutation;
use crate::admin::router::{AdminOperation, Operation, S3Router};
use crate::admin::runtime_sources::{
app_context_from_req, current_object_store_handle_for_context, current_scanner_metrics_report,
};
use crate::admin::runtime_sources::current_scanner_metrics_report;
use crate::module_switches::{ENV_SCANNER_ENABLED, scanner_enabled_from_env};
use crate::server::ADMIN_PREFIX;
use chrono::Utc;
@@ -25,13 +22,11 @@ use http::{HeaderMap, HeaderValue};
use hyper::{Method, StatusCode};
use matchit::Params;
use rustfs_common::metrics::{ScannerLifecycleExpirySnapshot, ScannerMaintenanceControlSnapshot, ScannerMetricsReport};
use rustfs_config::MAX_ADMIN_REQUEST_BODY_SIZE;
use rustfs_credentials::Credentials;
use rustfs_policy::policy::action::{Action, AdminAction};
use s3s::header::CONTENT_TYPE;
use s3s::{Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error};
use serde::{Deserialize, Serialize};
use tokio_util::sync::CancellationToken;
use serde::Serialize;
const JSON_CONTENT_TYPE: &str = "application/json";
@@ -43,13 +38,6 @@ struct ScannerStatusResponse {
metrics: ScannerMetricsReport,
cycle_schedule: rustfs_scanner::ScannerCycleScheduleStatus,
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)]
@@ -129,7 +117,6 @@ fn scanner_status_response(
metrics,
cycle_schedule,
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(),
AdminOperation(&ScannerStatusHandler {}),
)?;
r.insert(
Method::POST,
format!("{ADMIN_PREFIX}/v3/scanner/cycle-state/reset").as_str(),
AdminOperation(&ScannerCycleStateResetHandler {}),
)?;
r.insert(
Method::GET,
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
}
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)>> {
let mut headers = HeaderMap::new();
let content_type = HeaderValue::from_str(JSON_CONTENT_TYPE)
@@ -217,37 +192,6 @@ impl Operation for ScannerStatusHandler {
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]
impl Operation for IlmExpiryStatusHandler {
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"));
}
#[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]
fn scanner_disabled_reason_reports_startup_env_key() {
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"]["clean_idle_backoff_enabled"], false);
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]
-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::Put, "/rustfs/admin/v3/config", CONFIG_UPDATE, RouteRiskLevel::High),
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(
HttpMethod::Get,
"/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);
}
#[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]
fn route_policy_uses_tier_actions_for_transition_routes() {
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::PUT, "/v3/config"),
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_sample(
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::PUT, &admin_path("/v3/config"));
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::POST, &admin_path("/v3/ilm/transition/run"));
assert_route(
@@ -1369,7 +1367,6 @@ fn test_admin_alias_paths_match_existing_admin_routes() {
(Method::GET, compat_admin_alias_path("/v3/config")),
(Method::PUT, compat_admin_alias_path("/v3/config")),
(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")),
] {
assert!(
+1 -1
View File
@@ -1021,7 +1021,7 @@ fn build_list_objects_v2_metadata_output(
object: Object {
key: Some(encode_list_objects_v2_value(&object.name, encoding_type)),
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)),
storage_class: Some(ObjectStorageClass::from(
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
}
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
/// distributed delete path instead of its usual typed missing-object error.
fn is_delete_objects_not_found(error: &EcstoreError) -> bool {
@@ -8458,6 +8409,8 @@ impl DefaultObjectUsecase {
object: ObjectToDelete,
versioned: bool,
version_suspended: bool,
size: i64,
existing: Option<ObjectInfo>,
}
// Phase 2 (bounded concurrency, backlog#929 / HP-8): collect the
@@ -8475,23 +8428,32 @@ impl DefaultObjectUsecase {
skip_stat,
} = prepared;
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 {
Ok(_) => {}
Err(err) if is_err_object_not_found(&err) || is_err_version_not_found(&err) => {}
Ok(res) => (res, false),
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)),
}
}
};
let size = goi.size;
if synthetic_version_id {
object.version_id = Some(Uuid::nil());
}
let existing = (!skip_stat && !source_missing).then_some(goi);
Ok::<_, ApiError>(AdmittedDelete {
idx,
object,
versioned: opts.versioned,
version_suspended: opts.version_suspended,
size,
existing,
})
}))
.buffered(DELETE_OBJECTS_PRE_STAT_CONCURRENCY)
@@ -8502,11 +8464,15 @@ impl DefaultObjectUsecase {
// per-key success/failure reporting is unchanged.
let mut object_to_delete = 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();
for admitted in admitted_deletes {
object_sizes.push(admitted.size);
object_to_delete_idx.push(admitted.idx);
object_versioning.push((admitted.versioned, admitted.version_suspended));
object_to_delete.push(admitted.object);
existing_object_infos.push(admitted.existing);
}
let cache_adapter = self.object_data_cache();
let cache_keys_before_delete = object_to_delete
@@ -8523,8 +8489,8 @@ impl DefaultObjectUsecase {
..Default::default()
};
apply_bucket_generation_guard(&req, &bucket, &mut storage_delete_opts)?;
let (dobjs, errs, accounting) = store
.delete_objects_with_tier_delete_journal_and_accounting(&bucket, object_to_delete.clone(), storage_delete_opts)
let (dobjs, errs) = store
.delete_objects_with_tier_delete_journal(&bucket, object_to_delete.clone(), storage_delete_opts)
.await;
let _manager = get_concurrency_manager();
@@ -8549,16 +8515,17 @@ impl DefaultObjectUsecase {
delete_results[didx].delete_object = Some(deleted_object.clone());
let (versioned, version_suspended) = object_versioning[i];
let creates_delete_marker = object_to_delete[i].version_id.is_none() && versioned && !version_suspended;
let committed_delete_marker = dobjs[i].delete_marker;
let delete_accounting = accounting.get(i).and_then(Option::as_ref);
let update = delete_memory_update(
creates_delete_marker,
committed_delete_marker,
delete_request_targets_current(object_to_delete[i].version_id),
delete_accounting.and_then(|value| value.size),
delete_accounting.is_some_and(|value| value.removed_current_object),
);
apply_delete_memory_update(&bucket, update).await;
if creates_delete_marker {
record_bucket_delete_marker_memory(&bucket).await;
} else {
let size = object_sizes[i].max(0) as u64;
record_bucket_object_delete_memory(
&bucket,
size,
existing_object_infos[i].is_some() && object_to_delete[i].version_id.is_none(),
)
.await;
}
}
Err(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;
}
// Fast in-memory update for immediate quota and admin usage consistency.
// Prefix/force deletes and synthetic directory entries do not carry one
// committed object identity; leave their cache delta to reconciliation.
let update = if force_delete || obj_info.name.is_empty() || synthetic_version_id {
None
// Fast in-memory update for immediate quota and admin usage consistency
if delete_creates_delete_marker(&opts) {
record_bucket_delete_marker_memory(&bucket).await;
} else {
// The storage commit returns this object's metadata while its
// generation lock is held. Never fall back to a pre-delete stat:
// an overwrite can commit between that stat and this delete.
delete_memory_update(
delete_creates_delete_marker(&opts),
obj_info.delete_marker,
opts.version_id.is_none(),
quota_object_size(&obj_info).ok(),
delete_removes_current_object(&opts),
)
};
apply_delete_memory_update(&bucket, update).await;
record_bucket_object_delete_memory(&bucket, obj_info.size.max(0) as u64, opts.version_id.is_none()).await;
}
if obj_info.name.is_empty() {
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));
}
#[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]
async fn execute_get_object_attributes_returns_internal_error_when_store_uninitialized() {
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,
};
#[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) {
crate::storage::storage_api::ecstore_data_usage::record_bucket_object_delete_memory(
bucket,
@@ -1238,10 +1233,7 @@ pub(crate) mod test {
pub(crate) use super::access::ReqInfo;
pub(crate) use super::options::VERSIONING_CONFIG_LOOKUPS;
pub(crate) use super::{bucket, ecfs, object_utils, runtime};
pub(crate) mod data_usage {
pub(crate) use super::super::data_usage::*;
}
pub(crate) use super::{bucket, data_usage, ecfs, object_utils, runtime};
pub(crate) use crate::storage::storage_api::test_consumer::{get_global_bucket_metadata_sys, set_bucket_metadata};
pub(crate) use crate::storage::storage_api::{
ECStore, Endpoint, Endpoints, PoolEndpoints, StorageObjectInfo, StorageObjectOptions, StoragePutObjReader,
+173
View File
@@ -0,0 +1,173 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use std::env;
use std::ffi::OsString;
use std::fs;
use std::path::PathBuf;
use std::time::Duration;
use super::{CredentialStore, IdentityStore};
pub const ENV_CONNECT_ENDPOINT: &str = "RUSTFS_CONNECT_ENDPOINT";
pub const ENV_CONNECT_ROOT_CA_FILE: &str = "RUSTFS_CONNECT_ROOT_CA_FILE";
pub const ENV_CONNECT_STATE_DIR: &str = "RUSTFS_CONNECT_STATE_DIR";
#[derive(Clone, Copy, Debug)]
pub struct HeartbeatSchedule {
pub cadence: Duration,
pub jitter: Duration,
pub timeout: Duration,
pub initial_backoff: Duration,
pub max_backoff: Duration,
}
impl Default for HeartbeatSchedule {
fn default() -> Self {
Self {
cadence: Duration::from_secs(30),
jitter: Duration::from_secs(3),
timeout: Duration::from_secs(5),
initial_backoff: Duration::from_secs(1),
max_backoff: Duration::from_secs(5 * 60),
}
}
}
#[derive(Clone, Debug)]
pub struct HeartbeatConfig {
pub endpoint: String,
pub root_ca_pem: Vec<u8>,
pub identity_store: IdentityStore,
pub credential_store: CredentialStore,
pub state_path: PathBuf,
pub schedule: HeartbeatSchedule,
}
impl HeartbeatConfig {
pub fn new(
endpoint: impl Into<String>,
root_ca_pem: impl Into<Vec<u8>>,
identity_store: IdentityStore,
credential_store: CredentialStore,
state_path: impl Into<PathBuf>,
) -> Self {
Self {
endpoint: endpoint.into(),
root_ca_pem: root_ca_pem.into(),
identity_store,
credential_store,
state_path: state_path.into(),
schedule: HeartbeatSchedule::default(),
}
}
pub fn from_env() -> Result<Option<Self>, HeartbeatConfigError> {
Self::from_env_values(
env::var_os(ENV_CONNECT_ENDPOINT),
env::var_os(ENV_CONNECT_ROOT_CA_FILE),
env::var_os(ENV_CONNECT_STATE_DIR),
)
}
fn from_env_values(
endpoint: Option<OsString>,
root_ca_file: Option<OsString>,
state_dir: Option<OsString>,
) -> Result<Option<Self>, HeartbeatConfigError> {
let configured = endpoint.is_some() || root_ca_file.is_some() || state_dir.is_some();
if !configured {
return Ok(None);
}
let (Some(endpoint), Some(root_ca_file), Some(state_dir)) = (endpoint, root_ca_file, state_dir) else {
return Err(HeartbeatConfigError::Partial);
};
let endpoint = endpoint.into_string().map_err(|_| HeartbeatConfigError::EndpointEncoding)?;
let root_ca_file = PathBuf::from(root_ca_file);
let state_dir = PathBuf::from(state_dir);
if endpoint.is_empty() || root_ca_file.as_os_str().is_empty() || state_dir.as_os_str().is_empty() {
return Err(HeartbeatConfigError::Partial);
}
let root_ca_pem = fs::read(&root_ca_file).map_err(|source| HeartbeatConfigError::RootCertificate {
path: root_ca_file,
source,
})?;
Ok(Some(Self::new(
endpoint,
root_ca_pem,
IdentityStore::new(state_dir.join("identity")),
CredentialStore::new(state_dir.join("credential")),
state_dir.join("heartbeat/state.json"),
)))
}
}
#[derive(Debug, thiserror::Error)]
pub enum HeartbeatConfigError {
#[error(
"Connect heartbeat configuration requires RUSTFS_CONNECT_ENDPOINT, RUSTFS_CONNECT_ROOT_CA_FILE, and RUSTFS_CONNECT_STATE_DIR"
)]
Partial,
#[error("RUSTFS_CONNECT_ENDPOINT is not valid UTF-8")]
EndpointEncoding,
#[error("failed to read the Connect root CA at {path}: {source}")]
RootCertificate {
path: PathBuf,
#[source]
source: std::io::Error,
},
}
#[cfg(test)]
mod tests {
use super::{HeartbeatConfig, HeartbeatConfigError};
use std::ffi::OsString;
#[test]
fn absent_environment_is_disabled_without_side_effects() {
assert!(
HeartbeatConfig::from_env_values(None, None, None)
.expect("absent config")
.is_none()
);
}
#[test]
fn partial_environment_is_rejected() {
assert!(matches!(
HeartbeatConfig::from_env_values(Some(OsString::from("https://connect.example/agent/")), None, None),
Err(HeartbeatConfigError::Partial)
));
}
#[test]
fn complete_environment_builds_the_durable_paths() {
let temp = tempfile::tempdir().expect("tempdir");
let root = temp.path().join("root.pem");
std::fs::write(&root, b"root certificate").expect("root CA");
let state = temp.path().join("state");
let config = HeartbeatConfig::from_env_values(
Some(OsString::from("https://connect.example/agent/")),
Some(root.into_os_string()),
Some(state.clone().into_os_string()),
)
.expect("complete config")
.expect("enabled config");
assert_eq!(config.endpoint, "https://connect.example/agent/");
assert_eq!(config.root_ca_pem, b"root certificate");
assert_eq!(config.state_path, state.join("heartbeat/state.json"));
assert!(!state.exists(), "parsing configuration must not create state");
}
}
+585
View File
@@ -0,0 +1,585 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use std::fs;
use std::io::{self, Write as _};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Duration;
use chrono::{DateTime, SecondsFormat, Utc};
use reqwest::{Client, StatusCode, Url, header};
use rustls::RootCertStore;
use rustls::pki_types::{CertificateDer, pem::PemObject as _};
use serde::{Deserialize, Serialize};
use uuid::Uuid;
use zeroize::Zeroizing;
use super::config::HeartbeatConfig;
use super::credential_store::{CredentialStoreError, DeviceCredential};
use super::identity::IdentityError;
use super::identity_store::StoreError;
use super::registration::{CredentialValidationError, validate_stored_credential};
const PROTOCOL_VERSION: &str = "v1";
const AGENT_VERSION: &str = concat!("rustfs-agent/", env!("CARGO_PKG_VERSION"));
const MAX_SEQUENCE: u64 = 9_007_199_254_740_991;
const MAX_RESPONSE_BYTES: usize = 64 * 1024;
#[cfg(unix)]
const FILE_MODE: u32 = 0o600;
static STAGING_SEQUENCE: AtomicU64 = AtomicU64::new(0);
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields, rename_all = "camelCase")]
pub struct CoarseNodeSummary {
total: u16,
healthy: u16,
degraded: u16,
}
impl CoarseNodeSummary {
pub fn new(total: u16, healthy: u16, degraded: u16) -> Result<Self, HeartbeatError> {
let summary = Self {
total,
healthy,
degraded,
};
if !summary.is_valid() {
return Err(HeartbeatError::NodeSummary);
}
Ok(summary)
}
fn is_valid(&self) -> bool {
self.total != 0
&& self.total <= 4096
&& self.healthy <= 4096
&& self.degraded <= 4096
&& self.healthy.saturating_add(self.degraded) <= self.total
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum HeartbeatStatus {
Starting,
Online { server_time: String },
BackingOff { delay: Duration },
AuthenticationStopped { status: u16, reason: Option<String> },
Failed { reason: String },
Stopped,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields, rename_all = "camelCase")]
pub(crate) struct PendingHeartbeat {
protocol_version: String,
request_id: String,
agent_version: String,
capabilities: [String; 1],
sequence: u64,
client_time: String,
coarse_node_summary: CoarseNodeSummary,
}
impl PendingHeartbeat {
fn is_valid(&self) -> bool {
self.protocol_version == PROTOCOL_VERSION
&& self.agent_version == AGENT_VERSION
&& self.capabilities[0] == "heartbeat"
&& self.sequence <= MAX_SEQUENCE
&& self.coarse_node_summary.is_valid()
&& is_exact_utc_seconds(&self.client_time)
&& Uuid::parse_str(&self.request_id)
.is_ok_and(|request_id| request_id.get_version_num() == 4 && request_id.to_string() == self.request_id)
}
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct HeartbeatResponse {
server_time: String,
accepted_version: String,
#[serde(default)]
capability_hints: Vec<String>,
}
pub(crate) enum Delivery {
Accepted { server_time: String },
Retry { retry_after: Option<Duration> },
AuthenticationStopped { status: u16, reason: Option<String> },
Rejected { status: u16, reason: Option<String> },
}
pub(crate) struct HeartbeatSender {
endpoint: Url,
root_store: RootCertStore,
roots: Vec<CertificateDer<'static>>,
config: HeartbeatConfig,
}
impl HeartbeatSender {
pub(crate) fn new(config: HeartbeatConfig) -> Result<Self, HeartbeatError> {
let mut endpoint = Url::parse(&config.endpoint).map_err(|_| HeartbeatError::Endpoint)?;
if endpoint.scheme() != "https"
|| endpoint.cannot_be_a_base()
|| !endpoint.username().is_empty()
|| endpoint.password().is_some()
|| endpoint.query().is_some()
|| endpoint.fragment().is_some()
{
return Err(HeartbeatError::Endpoint);
}
if !endpoint.path().ends_with('/') {
endpoint.set_path(&format!("{}/", endpoint.path()));
}
let roots = CertificateDer::pem_slice_iter(&config.root_ca_pem)
.collect::<Result<Vec<_>, _>>()
.map_err(|_| HeartbeatError::RootCertificate)?;
if roots.is_empty() {
return Err(HeartbeatError::RootCertificate);
}
let mut root_store = RootCertStore::empty();
let (accepted, rejected) = root_store.add_parsable_certificates(roots.clone());
if accepted != roots.len() || rejected != 0 {
return Err(HeartbeatError::RootCertificate);
}
let schedule = config.schedule;
if schedule.cadence.is_zero()
|| schedule.timeout.is_zero()
|| schedule.timeout > Duration::from_secs(5)
|| schedule.initial_backoff.is_zero()
|| schedule.max_backoff < schedule.initial_backoff
|| schedule.max_backoff > Duration::from_secs(5 * 60)
|| schedule.jitter > schedule.cadence
{
return Err(HeartbeatError::Schedule);
}
Ok(Self {
endpoint,
root_store,
roots,
config,
})
}
pub(crate) async fn send(&self, heartbeat: &PendingHeartbeat) -> Result<Delivery, HeartbeatError> {
let (cluster_uid, client) = {
let _lock = self.config.credential_store.lock().await?;
let credential = self.config.credential_store.load()?.ok_or(HeartbeatError::NotRegistered)?;
let identity = self.config.identity_store.load()?.ok_or(HeartbeatError::IdentityMissing)?;
validate_stored_credential(&credential, &identity, &self.root_store, &self.roots)?;
let now = Utc::now().timestamp();
if now < credential.not_before_unix || now >= credential.not_after_unix {
return Err(HeartbeatError::CredentialExpired);
}
let cluster_uid = cluster_uid(&credential)?.to_owned();
let client = self.client(&credential, &identity.to_pkcs8_pem()?)?;
(cluster_uid, client)
};
let url = self.endpoint.join(&format!("clusters/{cluster_uid}/heartbeats"))?;
let response = match client.post(url).json(heartbeat).send().await {
Ok(response) => response,
Err(error) if error.is_timeout() || error.is_connect() || error.is_request() => {
return Ok(Delivery::Retry { retry_after: None });
}
Err(error) => return Err(error.into()),
};
let status = response.status();
if status == StatusCode::TOO_MANY_REQUESTS {
return Ok(Delivery::Retry {
retry_after: retry_after(response.headers(), Utc::now(), self.config.schedule.max_backoff),
});
}
if status == StatusCode::REQUEST_TIMEOUT || status.is_server_error() {
return Ok(Delivery::Retry { retry_after: None });
}
if matches!(status, StatusCode::UNAUTHORIZED | StatusCode::FORBIDDEN) {
return Ok(Delivery::AuthenticationStopped {
status: status.as_u16(),
reason: response_reason(response).await,
});
}
if status != StatusCode::OK {
return Ok(Delivery::Rejected {
status: status.as_u16(),
reason: response_reason(response).await,
});
}
let accepted: HeartbeatResponse =
serde_json::from_slice(&bounded_body(response).await?).map_err(|_| HeartbeatError::Response)?;
if accepted.accepted_version != PROTOCOL_VERSION
|| accepted.capability_hints.len() > 32
|| accepted.capability_hints.iter().any(|hint| hint.len() > 32)
|| !is_exact_utc_seconds(&accepted.server_time)
{
return Err(HeartbeatError::Response);
}
Ok(Delivery::Accepted {
server_time: accepted.server_time,
})
}
fn client(&self, credential: &DeviceCredential, key: &Zeroizing<String>) -> Result<Client, HeartbeatError> {
let mut pem = Zeroizing::new(Vec::with_capacity(credential.certificate_chain.len() + key.len() + 1));
pem.extend_from_slice(credential.certificate_chain.as_bytes());
pem.push(b'\n');
pem.extend_from_slice(key.as_bytes());
let identity = reqwest::Identity::from_pem(&pem).map_err(|_| HeartbeatError::IdentityCertificate)?;
let roots = self
.roots
.iter()
.map(|root| reqwest::Certificate::from_der(root.as_ref()))
.collect::<Result<Vec<_>, _>>()?;
Client::builder()
.https_only(true)
.redirect(reqwest::redirect::Policy::none())
.timeout(self.config.schedule.timeout)
.tls_certs_only(roots)
.identity(identity)
.build()
.map_err(Into::into)
}
}
#[derive(Clone)]
pub(crate) struct HeartbeatStateStore {
path: PathBuf,
}
#[derive(Default, Serialize, Deserialize)]
#[serde(deny_unknown_fields, rename_all = "camelCase")]
struct HeartbeatState {
next_sequence: u64,
pending: Option<PendingHeartbeat>,
}
impl HeartbeatStateStore {
pub(crate) fn new(path: PathBuf) -> Self {
Self { path }
}
pub(crate) fn try_runtime_lock(&self) -> Result<fs::File, HeartbeatError> {
let directory = parent(&self.path)?;
fs::create_dir_all(directory).map_err(|source| state_io(directory, source))?;
let name = filename(&self.path)?;
let path = directory.join(format!(".{name}.lock"));
let mut options = fs::OpenOptions::new();
options.create(true).truncate(false).read(true).write(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt as _;
options.mode(FILE_MODE);
}
let lock = options.open(&path).map_err(|source| state_io(&path, source))?;
check_mode(&path)?;
lock.try_lock().map_err(|_| HeartbeatError::AlreadyRunning)?;
Ok(lock)
}
pub(crate) async fn prepare(
&self,
summary: CoarseNodeSummary,
now: DateTime<Utc>,
) -> Result<PendingHeartbeat, HeartbeatError> {
let store = self.clone();
tokio::task::spawn_blocking(move || store.prepare_sync(summary, now))
.await
.map_err(|source| state_io(&self.path, io::Error::other(source)))?
}
pub(crate) async fn mark_accepted(&self, accepted: &PendingHeartbeat) -> Result<(), HeartbeatError> {
let store = self.clone();
let accepted = accepted.clone();
tokio::task::spawn_blocking(move || store.mark_accepted_sync(&accepted))
.await
.map_err(|source| state_io(&self.path, io::Error::other(source)))?
}
fn prepare_sync(&self, summary: CoarseNodeSummary, now: DateTime<Utc>) -> Result<PendingHeartbeat, HeartbeatError> {
let mut state = self.read()?;
if let Some(pending) = state.pending {
return Ok(pending);
}
if state.next_sequence > MAX_SEQUENCE {
return Err(HeartbeatError::SequenceExhausted);
}
let pending = PendingHeartbeat {
protocol_version: PROTOCOL_VERSION.to_owned(),
request_id: Uuid::new_v4().to_string(),
agent_version: AGENT_VERSION.to_owned(),
capabilities: ["heartbeat".to_owned()],
sequence: state.next_sequence,
client_time: now.to_rfc3339_opts(SecondsFormat::Secs, true),
coarse_node_summary: summary,
};
state.pending = Some(pending.clone());
self.write(&state)?;
Ok(pending)
}
fn mark_accepted_sync(&self, accepted: &PendingHeartbeat) -> Result<(), HeartbeatError> {
let mut state = self.read()?;
if state.pending.as_ref() != Some(accepted) {
return Err(HeartbeatError::StateConflict);
}
state.next_sequence = accepted.sequence.checked_add(1).ok_or(HeartbeatError::SequenceExhausted)?;
state.pending = None;
self.write(&state)
}
fn read(&self) -> Result<HeartbeatState, HeartbeatError> {
let bytes = match fs::read(&self.path) {
Ok(bytes) => bytes,
Err(source) if source.kind() == io::ErrorKind::NotFound => return Ok(HeartbeatState::default()),
Err(source) => return Err(state_io(&self.path, source)),
};
check_mode(&self.path)?;
let state: HeartbeatState = serde_json::from_slice(&bytes).map_err(|source| HeartbeatError::StateInvalid {
path: self.path.clone(),
source,
})?;
if state.next_sequence > MAX_SEQUENCE + 1
|| state
.pending
.as_ref()
.is_some_and(|pending| pending.sequence != state.next_sequence || !pending.is_valid())
{
return Err(HeartbeatError::StateCorrupt { path: self.path.clone() });
}
Ok(state)
}
fn write(&self, state: &HeartbeatState) -> Result<(), HeartbeatError> {
let bytes = serde_json::to_vec(state).map_err(|source| HeartbeatError::StateInvalid {
path: self.path.clone(),
source,
})?;
let directory = parent(&self.path)?;
fs::create_dir_all(directory).map_err(|source| state_io(directory, source))?;
let temp = stage(directory, &self.path, &bytes)?;
let result = fs::rename(&temp, &self.path)
.map_err(|source| state_io(&self.path, source))
.and_then(|()| fsync_dir(directory).map_err(|source| state_io(directory, source)));
if result.is_err() {
let _ = fs::remove_file(temp);
}
result
}
}
fn cluster_uid(credential: &DeviceCredential) -> Result<&str, HeartbeatError> {
let mut parts = credential.name.split('/');
let valid = parts.next() == Some("organizations");
let organization_uid = parts.next();
let valid = valid && parts.next() == Some("clusters");
let cluster_uid = parts.next();
let valid = valid && parts.next() == Some("clusterDevices");
let device_uid = parts.next();
if !valid
|| organization_uid.is_none_or(str::is_empty)
|| cluster_uid.is_none_or(str::is_empty)
|| device_uid != Some(credential.uid.as_str())
|| parts.next().is_some()
{
return Err(HeartbeatError::CredentialName);
}
cluster_uid.ok_or(HeartbeatError::CredentialName)
}
fn retry_after(headers: &header::HeaderMap, now: DateTime<Utc>, maximum: Duration) -> Option<Duration> {
let value = headers.get(header::RETRY_AFTER)?.to_str().ok()?;
let delay = value.parse::<u64>().ok().map(Duration::from_secs).or_else(|| {
DateTime::parse_from_rfc2822(value)
.ok()
.and_then(|at| (at.with_timezone(&Utc) - now).to_std().ok())
})?;
Some(delay.min(maximum))
}
fn is_exact_utc_seconds(value: &str) -> bool {
DateTime::parse_from_rfc3339(value).is_ok_and(|time| {
time.offset().local_minus_utc() == 0
&& value.ends_with('Z')
&& time.with_timezone(&Utc).to_rfc3339_opts(SecondsFormat::Secs, true) == value
})
}
async fn response_reason(response: reqwest::Response) -> Option<String> {
#[derive(Deserialize)]
struct Envelope {
#[serde(default)]
details: Vec<Detail>,
}
#[derive(Deserialize)]
struct Detail {
#[serde(default)]
reason: String,
}
serde_json::from_slice::<Envelope>(&bounded_body(response).await.ok()?)
.ok()?
.details
.into_iter()
.find_map(|detail| (!detail.reason.is_empty()).then_some(detail.reason))
}
async fn bounded_body(mut response: reqwest::Response) -> Result<Vec<u8>, HeartbeatError> {
let mut body = Vec::new();
while let Some(chunk) = response.chunk().await? {
if body.len().saturating_add(chunk.len()) > MAX_RESPONSE_BYTES {
return Err(HeartbeatError::ResponseTooLarge);
}
body.extend_from_slice(&chunk);
}
Ok(body)
}
fn parent(path: &Path) -> Result<&Path, HeartbeatError> {
path.parent()
.ok_or_else(|| state_io(path, io::Error::new(io::ErrorKind::InvalidInput, "state path has no parent")))
}
fn filename(path: &Path) -> Result<&str, HeartbeatError> {
path.file_name()
.and_then(|name| name.to_str())
.ok_or_else(|| state_io(path, io::Error::new(io::ErrorKind::InvalidInput, "state filename is invalid")))
}
fn stage(directory: &Path, destination: &Path, bytes: &[u8]) -> Result<PathBuf, HeartbeatError> {
let name = filename(destination)?;
loop {
let path = directory.join(format!(
".{name}.{}.{}.tmp",
std::process::id(),
STAGING_SEQUENCE.fetch_add(1, Ordering::Relaxed)
));
let mut options = fs::OpenOptions::new();
options.write(true).create_new(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt as _;
options.mode(FILE_MODE);
}
let mut file = match options.open(&path) {
Ok(file) => file,
Err(source) if source.kind() == io::ErrorKind::AlreadyExists => continue,
Err(source) => return Err(state_io(&path, source)),
};
if let Err(source) = file.write_all(bytes).and_then(|()| file.sync_all()) {
let _ = fs::remove_file(&path);
return Err(state_io(&path, source));
}
return Ok(path);
}
}
fn state_io(path: &Path, source: io::Error) -> HeartbeatError {
HeartbeatError::StateIo {
path: path.to_path_buf(),
source,
}
}
#[cfg(unix)]
fn check_mode(path: &Path) -> Result<(), HeartbeatError> {
use std::os::unix::fs::PermissionsExt as _;
let mode = fs::metadata(path)
.map_err(|source| state_io(path, source))?
.permissions()
.mode()
& 0o7777;
if mode != FILE_MODE {
return Err(HeartbeatError::StatePermissions {
path: path.to_path_buf(),
mode,
expected: FILE_MODE,
});
}
Ok(())
}
#[cfg(not(unix))]
fn check_mode(_path: &Path) -> Result<(), HeartbeatError> {
Ok(())
}
fn fsync_dir(directory: &Path) -> io::Result<()> {
#[cfg(unix)]
fs::File::open(directory)?.sync_all()?;
#[cfg(not(unix))]
let _ = directory;
Ok(())
}
#[derive(Debug, thiserror::Error)]
pub enum HeartbeatError {
#[error("Connect heartbeat endpoint must be an HTTPS base URL without credentials, query, or fragment")]
Endpoint,
#[error("Connect heartbeat root CA configuration is invalid")]
RootCertificate,
#[error("Connect heartbeat schedule is invalid")]
Schedule,
#[error("RustFS is not registered with Connect")]
NotRegistered,
#[error("the Connect device private key is missing")]
IdentityMissing,
#[error("the stored Connect certificate and device private key cannot form a TLS identity")]
IdentityCertificate,
#[error("the stored Connect credential name is invalid")]
CredentialName,
#[error("the stored Connect device certificate is not currently valid")]
CredentialExpired,
#[error("the Connect heartbeat node summary is outside protocol bounds")]
NodeSummary,
#[error("the Connect heartbeat sequence is exhausted")]
SequenceExhausted,
#[error("a Connect heartbeat runtime already owns this state")]
AlreadyRunning,
#[error("the persisted Connect heartbeat changed while delivery was in flight")]
StateConflict,
#[error("Connect heartbeat state I/O failed at {path}: {source}")]
StateIo {
path: PathBuf,
#[source]
source: io::Error,
},
#[error("Connect heartbeat state at {path} is invalid: {source}")]
StateInvalid {
path: PathBuf,
#[source]
source: serde_json::Error,
},
#[error("Connect heartbeat state at {path} violates the protocol invariants")]
StateCorrupt { path: PathBuf },
#[cfg(unix)]
#[error("Connect heartbeat state at {path} has mode {mode:o}, expected {expected:o}")]
StatePermissions { path: PathBuf, mode: u32, expected: u32 },
#[error("Connect heartbeat response exceeded 64 KiB")]
ResponseTooLarge,
#[error("Connect returned an invalid heartbeat response")]
Response,
#[error(transparent)]
Url(#[from] url::ParseError),
#[error(transparent)]
Transport(#[from] reqwest::Error),
#[error(transparent)]
Identity(#[from] IdentityError),
#[error(transparent)]
IdentityStore(#[from] StoreError),
#[error(transparent)]
CredentialStore(#[from] CredentialStoreError),
#[error(transparent)]
CredentialValidation(#[from] CredentialValidationError),
}
+9 -3
View File
@@ -21,20 +21,26 @@
//! canonical transcript frozen by
//! `protocol/agent/v1/registration-proof.md`.
//!
//! Nothing here contacts the network or starts a task. A deployment that has
//! not been enrolled into a Connect control plane never calls into it, so an
//! unconfigured server generates no key and holds no identity.
//! Enrolled deployments may start the optional outbound heartbeat runtime.
//! An unconfigured server starts no Connect task, generates no key, and holds
//! no Connect identity.
pub mod client;
pub mod config;
pub mod credential_store;
pub mod heartbeat;
pub mod identity;
pub mod identity_store;
pub mod offline;
pub mod registration;
pub mod runtime;
pub use client::{ClientError, ConnectClient, ConnectConfig};
pub use config::{HeartbeatConfig, HeartbeatConfigError, HeartbeatSchedule};
pub use credential_store::{CredentialStore, DeviceCredential};
pub use heartbeat::{CoarseNodeSummary, HeartbeatError, HeartbeatStatus};
pub use identity::{DeviceIdentity, IdentityError, RegistrationProof, RegistrationTranscript};
pub use identity_store::{IdentityStore, StoreError};
pub use offline::{EnrollmentError, OfflineEnrollment, OfflineKeyStore, VerifiedChallenge};
pub use registration::{RegistrationToken, TokenError};
pub use runtime::{HeartbeatRuntime, spawn_heartbeat_runtime};
+156
View File
@@ -0,0 +1,156 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use std::future::Future;
use std::time::Duration;
use chrono::Utc;
use rand::RngExt as _;
use tokio::sync::watch;
use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
use super::config::HeartbeatConfig;
use super::heartbeat::{CoarseNodeSummary, Delivery, HeartbeatError, HeartbeatSender, HeartbeatStateStore, HeartbeatStatus};
pub struct HeartbeatRuntime {
shutdown: CancellationToken,
status: watch::Receiver<HeartbeatStatus>,
task: Option<JoinHandle<()>>,
}
impl HeartbeatRuntime {
pub fn status(&self) -> watch::Receiver<HeartbeatStatus> {
self.status.clone()
}
pub async fn shutdown(mut self) {
self.shutdown.cancel();
if let Some(task) = self.task.take() {
let _ = task.await;
}
}
}
impl Drop for HeartbeatRuntime {
fn drop(&mut self) {
self.shutdown.cancel();
}
}
pub fn spawn_heartbeat_runtime<F>(
config: Option<HeartbeatConfig>,
parent_shutdown: &CancellationToken,
sample: F,
) -> Result<Option<HeartbeatRuntime>, HeartbeatError>
where
F: Fn() -> CoarseNodeSummary + Send + Sync + 'static,
{
let Some(config) = config else {
return Ok(None);
};
let sender = HeartbeatSender::new(config.clone())?;
let store = HeartbeatStateStore::new(config.state_path.clone());
let lock = store.try_runtime_lock()?;
let schedule = config.schedule;
let shutdown = parent_shutdown.child_token();
let task_shutdown = shutdown.clone();
let (status_tx, status_rx) = watch::channel(HeartbeatStatus::Starting);
let task = tokio::spawn(async move {
let _lock = lock;
let mut backoff = schedule.initial_backoff;
loop {
if task_shutdown.is_cancelled() {
break;
}
let pending = match store.prepare(sample(), Utc::now()).await {
Ok(pending) => pending,
Err(error) => return failed(&status_tx, error),
};
let delivery = match cancellable(&task_shutdown, sender.send(&pending)).await {
Some(Ok(delivery)) => delivery,
Some(Err(error)) => return failed(&status_tx, error),
None => break,
};
let delay = match delivery {
Delivery::Accepted { server_time } => {
if let Err(error) = store.mark_accepted(&pending).await {
return failed(&status_tx, error);
}
backoff = schedule.initial_backoff;
let _ = status_tx.send(HeartbeatStatus::Online { server_time });
schedule.cadence.saturating_add(jitter(schedule.jitter))
}
Delivery::Retry { retry_after } => {
let delay = retry_after
.unwrap_or(backoff)
.clamp(schedule.initial_backoff, schedule.max_backoff);
backoff = backoff.saturating_mul(2).min(schedule.max_backoff);
let _ = status_tx.send(HeartbeatStatus::BackingOff { delay });
delay
}
Delivery::AuthenticationStopped { status, reason } => {
let _ = status_tx.send(HeartbeatStatus::AuthenticationStopped { status, reason });
return;
}
Delivery::Rejected { status, reason } => {
let suffix = reason.map_or_else(String::new, |reason| format!("; reason={reason}"));
let _ = status_tx.send(HeartbeatStatus::Failed {
reason: format!("Connect rejected heartbeat with HTTP {status}{suffix}"),
});
return;
}
};
if sleep_or_cancel(&task_shutdown, delay).await {
break;
}
}
let _ = status_tx.send(HeartbeatStatus::Stopped);
});
Ok(Some(HeartbeatRuntime {
shutdown,
status: status_rx,
task: Some(task),
}))
}
fn failed(status: &watch::Sender<HeartbeatStatus>, error: HeartbeatError) {
let _ = status.send(HeartbeatStatus::Failed {
reason: error.to_string(),
});
}
fn jitter(maximum: Duration) -> Duration {
if maximum.is_zero() {
Duration::ZERO
} else {
maximum.mul_f64(rand::rng().random_range(0.0..=1.0))
}
}
async fn cancellable<T>(shutdown: &CancellationToken, future: impl Future<Output = T>) -> Option<T> {
tokio::select! {
biased;
() = shutdown.cancelled() => None,
value = future => Some(value),
}
}
async fn sleep_or_cancel(shutdown: &CancellationToken, delay: Duration) -> bool {
tokio::select! {
biased;
() = shutdown.cancelled() => true,
() = tokio::time::sleep(delay) => false,
}
}
+4
View File
@@ -128,6 +128,7 @@ pub(crate) async fn run_startup_runtime_lifecycle(lifecycle: StartupRuntimeLifec
} = lifecycle;
let StartupServiceRuntime {
optional_runtimes,
heartbeat,
iam_bootstrap,
enable_scanner,
} = service_runtime;
@@ -162,6 +163,9 @@ pub(crate) async fn run_startup_runtime_lifecycle(lifecycle: StartupRuntimeLifec
shutdown_token,
)
.await;
if let Some(heartbeat) = heartbeat {
heartbeat.shutdown().await;
}
if let Err(err) = event_notifier_reconciler.await {
tracing::warn!(
target: "rustfs::main::run",
+21
View File
@@ -16,6 +16,7 @@ use crate::site_replication_reconcile::spawn_site_replication_reconcile_task;
use crate::storage_api::startup::services::{ECStore, EndpointServerPools, ServerContextSlot};
use crate::{
config::Config,
connect::{CoarseNodeSummary, HeartbeatConfig, HeartbeatRuntime, spawn_heartbeat_runtime},
init::{init_buffer_profile_system, init_kms_system},
server::ServiceStateManager,
startup_audit::init_audit_runtime,
@@ -35,6 +36,7 @@ use tokio_util::sync::CancellationToken;
pub(crate) struct StartupServiceRuntime {
pub(crate) optional_runtimes: OptionalRuntimeServices,
pub(crate) heartbeat: Option<HeartbeatRuntime>,
pub(crate) iam_bootstrap: IamBootstrapDisposition,
pub(crate) enable_scanner: bool,
}
@@ -73,6 +75,8 @@ pub(crate) async fn init_startup_runtime_services(
init_kms_system(config).await?;
let optional_runtimes = init_optional_runtime_services().await?;
let heartbeat_config = HeartbeatConfig::from_env().map_err(std::io::Error::other)?;
let heartbeat_nodes = heartbeat_config.as_ref().map(|_| endpoint_pools.get_nodes().len());
init_buffer_profile_system(config);
init_deadlock_detector_runtime();
@@ -92,10 +96,27 @@ pub(crate) async fn init_startup_runtime_services(
init_notification_runtime(endpoint_pools, buckets).await?;
let enable_scanner = init_background_service_runtime(store.clone()).await?;
init_observability_runtime(store.clone(), ctx.clone()).await;
let heartbeat = start_heartbeat_runtime(heartbeat_config, heartbeat_nodes, &ctx)?;
Ok(StartupServiceRuntime {
optional_runtimes,
heartbeat,
iam_bootstrap,
enable_scanner,
})
}
fn start_heartbeat_runtime(
config: Option<HeartbeatConfig>,
node_count: Option<usize>,
shutdown: &CancellationToken,
) -> Result<Option<HeartbeatRuntime>> {
let Some(config) = config else {
return Ok(None);
};
let summary = u16::try_from(node_count.unwrap_or_default())
.ok()
.and_then(|total| CoarseNodeSummary::new(total, 0, 0).ok())
.ok_or_else(|| std::io::Error::other("Connect heartbeat node count is outside protocol bounds"))?;
spawn_heartbeat_runtime(Some(config), shutdown, move || summary).map_err(std::io::Error::other)
}
+1 -43
View File
@@ -296,10 +296,7 @@ pub(crate) fn build_list_objects_v2_output(
let mut obj = Object {
key: Some(key),
last_modified: v.mod_time.map(Timestamp::from),
// Compressed legacy objects may retain an unknown (-1)
// logical-size sentinel; never expose that internal value in
// an S3 response.
size: Some(v.get_actual_size_or_physical()),
size: Some(v.get_actual_size().unwrap_or_default()),
e_tag: v.etag.clone().map(|etag| to_s3s_etag(&etag)),
storage_class: v.storage_class.clone().map(ObjectStorageClass::from),
..Default::default()
@@ -659,45 +656,6 @@ mod tests {
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]
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");
-2
View File
@@ -429,8 +429,6 @@ pub(crate) mod ecstore_config {
}
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::{
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,
+687
View File
@@ -0,0 +1,687 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use std::collections::VecDeque;
use std::fs;
use std::path::Path;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use bytes::Bytes;
use http_body_util::{BodyExt as _, Full};
use hyper::service::service_fn;
use hyper::{Request, Response, StatusCode};
use hyper_util::rt::TokioIo;
use rcgen::{
BasicConstraints, CertificateParams, DistinguishedName, DnType, ExtendedKeyUsagePurpose, IsCa, Issuer, KeyPair,
KeyUsagePurpose, SanType,
};
use rustfs::connect::{
CoarseNodeSummary, CredentialStore, DeviceCredential, HeartbeatConfig, HeartbeatSchedule, HeartbeatStatus, IdentityStore,
spawn_heartbeat_runtime,
};
use rustls::RootCertStore;
use rustls::pki_types::{CertificateDer, PrivateKeyDer, PrivatePkcs8KeyDer};
use rustls::server::WebPkiClientVerifier;
use serde_json::{Value, json};
use time::OffsetDateTime;
use tokio::net::TcpListener;
use tokio::sync::watch;
use tokio_rustls::TlsAcceptor;
use tokio_util::sync::CancellationToken;
const ORGANIZATION_UID: &str = "0198f4b0-1a00-7c10-8d21-2e3f4a5b6c70";
const CLUSTER_UID: &str = "0198f4b0-2b00-7d20-9e31-3f4a5b6c7d81";
const DEVICE_UID: &str = "0198f4b0-3c00-7e30-8f41-4a5b6c7d8e92";
struct TestPki {
root_params: CertificateParams,
root_key: KeyPair,
root_der: CertificateDer<'static>,
root_pem: String,
server_der: CertificateDer<'static>,
server_key: PrivatePkcs8KeyDer<'static>,
}
impl TestPki {
fn new() -> Self {
let now = OffsetDateTime::now_utc();
let root_key = KeyPair::generate().expect("generate root key");
let mut root_params = CertificateParams::default();
root_params.is_ca = IsCa::Ca(BasicConstraints::Unconstrained);
root_params.not_before = now - time::Duration::days(30);
root_params.not_after = now + time::Duration::days(30);
root_params.key_usages = vec![KeyUsagePurpose::KeyCertSign, KeyUsagePurpose::DigitalSignature];
let root = root_params.self_signed(&root_key).expect("sign root");
let server_key = KeyPair::generate().expect("generate server key");
let mut server_params = CertificateParams::default();
server_params.not_before = now - time::Duration::hours(1);
server_params.not_after = now + time::Duration::days(2);
server_params
.subject_alt_names
.push(SanType::DnsName("localhost".try_into().expect("valid DNS name")));
server_params.extended_key_usages = vec![ExtendedKeyUsagePurpose::ServerAuth];
let server = server_params
.signed_by(&server_key, &Issuer::from_params(&root_params, &root_key))
.expect("sign server certificate");
Self {
root_params,
root_key,
root_der: root.der().clone(),
root_pem: root.pem(),
server_der: server.der().clone(),
server_key: PrivatePkcs8KeyDer::from(server_key.serialize_der()),
}
}
fn server_config(&self) -> rustls::ServerConfig {
let mut roots = RootCertStore::empty();
roots.add(self.root_der.clone()).expect("add client root");
let verifier = WebPkiClientVerifier::builder(Arc::new(roots))
.build()
.expect("client verifier");
rustls::ServerConfig::builder()
.with_client_cert_verifier(verifier)
.with_single_cert(vec![self.server_der.clone()], PrivateKeyDer::Pkcs8(self.server_key.clone_key()))
.expect("server TLS")
}
fn stores(&self, temp: &tempfile::TempDir) -> (IdentityStore, CredentialStore) {
let now = OffsetDateTime::now_utc();
self.stores_with_certificate(temp, now - time::Duration::hours(1), now + time::Duration::hours(23), true)
}
fn stores_with_certificate(
&self,
temp: &tempfile::TempDir,
not_before: OffsetDateTime,
not_after: OffsetDateTime,
bind_identity: bool,
) -> (IdentityStore, CredentialStore) {
let identity_store = IdentityStore::new(temp.path().join("identity"));
let identity = identity_store.load_or_create().expect("create identity");
let private_key = PrivatePkcs8KeyDer::from(identity.to_pkcs8_der().expect("serialize key").to_vec());
let device_key = if bind_identity {
KeyPair::from_pkcs8_der_and_sign_algo(&private_key, &rcgen::PKCS_ECDSA_P256_SHA256).expect("device key")
} else {
KeyPair::generate().expect("mismatched device key")
};
let mut params = CertificateParams::default();
params.not_before = not_before;
params.not_after = not_after;
params.serial_number = Some(vec![1; 16].into());
params.key_usages = vec![KeyUsagePurpose::DigitalSignature];
params.extended_key_usages = vec![ExtendedKeyUsagePurpose::ClientAuth];
params.distinguished_name = DistinguishedName::new();
params.distinguished_name.push(DnType::CommonName, DEVICE_UID);
params.subject_alt_names.push(SanType::URI(
format!("urn:rustfs:connect:device:{DEVICE_UID}")
.try_into()
.expect("device URI"),
));
let certificate = params
.signed_by(&device_key, &Issuer::from_params(&self.root_params, &self.root_key))
.expect("device certificate");
let cluster = format!("organizations/{ORGANIZATION_UID}/clusters/{CLUSTER_UID}");
let credential = DeviceCredential {
name: format!("{cluster}/clusterDevices/{DEVICE_UID}"),
uid: DEVICE_UID.to_owned(),
protocol_version: "v1".to_owned(),
key_id: format!("x509-{}", "01".repeat(16)),
certificate_serial: "01".repeat(16),
certificate: certificate.pem(),
certificate_chain: certificate.pem(),
not_before_unix: not_before.unix_timestamp(),
not_after_unix: not_after.unix_timestamp(),
};
let directory = temp.path().join("credential");
fs::create_dir_all(&directory).expect("credential directory");
let path = directory.join("device.crt.json");
fs::write(&path, serde_json::to_vec(&credential).expect("credential JSON")).expect("write credential");
private_mode(&path);
(identity_store, CredentialStore::new(directory))
}
}
#[derive(Clone)]
struct Reply {
status: StatusCode,
body: Value,
retry_after: Option<&'static str>,
delay: Duration,
}
impl Reply {
fn ok(time: &str) -> Self {
Self {
status: StatusCode::OK,
body: json!({
"serverTime": time,
"acceptedVersion": "v1",
"capabilityHints": [],
"futureField": true
}),
retry_after: None,
delay: Duration::ZERO,
}
}
fn error(status: StatusCode) -> Self {
Self {
status,
body: json!({"details": []}),
retry_after: None,
delay: Duration::ZERO,
}
}
}
struct TestServer {
endpoint: String,
seen: Arc<Mutex<Vec<Value>>>,
task: tokio::task::JoinHandle<()>,
}
impl Drop for TestServer {
fn drop(&mut self) {
self.task.abort();
}
}
async fn server(pki: &TestPki, replies: Vec<Reply>) -> TestServer {
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind server");
let address = listener.local_addr().expect("server address");
let acceptor = TlsAcceptor::from(Arc::new(pki.server_config()));
let replies = Arc::new(Mutex::new(VecDeque::from(replies)));
let seen = Arc::new(Mutex::new(Vec::new()));
let captured = seen.clone();
let task = tokio::spawn(async move {
while let Ok((stream, _)) = listener.accept().await {
let acceptor = acceptor.clone();
let replies = replies.clone();
let seen = captured.clone();
tokio::spawn(async move {
let Ok(stream) = acceptor.accept(stream).await else { return };
let service = service_fn(move |request: Request<hyper::body::Incoming>| {
let replies = replies.clone();
let seen = seen.clone();
async move {
assert_eq!(request.uri().path(), format!("/agent/clusters/{CLUSTER_UID}/heartbeats"));
let body = request.into_body().collect().await.expect("request body").to_bytes();
seen.lock()
.expect("seen lock")
.push(serde_json::from_slice(&body).expect("request JSON"));
let reply = replies
.lock()
.expect("reply lock")
.pop_front()
.unwrap_or_else(|| Reply::error(StatusCode::SERVICE_UNAVAILABLE));
if !reply.delay.is_zero() {
tokio::time::sleep(reply.delay).await;
}
let mut builder = Response::builder()
.status(reply.status)
.header("content-type", "application/json");
if let Some(value) = reply.retry_after {
builder = builder.header("retry-after", value);
}
Ok::<_, hyper::Error>(
builder
.body(Full::new(Bytes::from(serde_json::to_vec(&reply.body).expect("reply JSON"))))
.expect("reply"),
)
}
});
let _ = hyper::server::conn::http1::Builder::new()
.serve_connection(TokioIo::new(stream), service)
.await;
});
}
});
TestServer {
endpoint: format!("https://localhost:{}/agent/", address.port()),
seen,
task,
}
}
fn config(temp: &tempfile::TempDir, pki: &TestPki, server: &TestServer) -> HeartbeatConfig {
let (identity_store, credential_store) = pki.stores(temp);
config_with_stores(temp, pki, server, identity_store, credential_store)
}
fn config_with_stores(
temp: &tempfile::TempDir,
pki: &TestPki,
server: &TestServer,
identity_store: IdentityStore,
credential_store: CredentialStore,
) -> HeartbeatConfig {
HeartbeatConfig {
endpoint: server.endpoint.clone(),
root_ca_pem: pki.root_pem.as_bytes().to_vec(),
identity_store,
credential_store,
state_path: temp.path().join("heartbeat/state.json"),
schedule: HeartbeatSchedule {
cadence: Duration::from_millis(40),
jitter: Duration::ZERO,
timeout: Duration::from_millis(200),
initial_backoff: Duration::from_millis(20),
max_backoff: Duration::from_millis(80),
},
}
}
fn rewrite_credential(temp: &tempfile::TempDir, update: impl FnOnce(&mut DeviceCredential)) {
let path = temp.path().join("credential/device.crt.json");
let mut credential: DeviceCredential =
serde_json::from_slice(&fs::read(&path).expect("read credential")).expect("parse credential");
update(&mut credential);
fs::write(&path, serde_json::to_vec(&credential).expect("credential JSON")).expect("rewrite credential");
private_mode(&path);
}
fn summary() -> CoarseNodeSummary {
CoarseNodeSummary::new(8, 7, 1).expect("node summary")
}
async fn wait_for(
status: &mut watch::Receiver<HeartbeatStatus>,
predicate: impl Fn(&HeartbeatStatus) -> bool,
) -> HeartbeatStatus {
tokio::time::timeout(Duration::from_secs(3), async {
loop {
let current = status.borrow_and_update().clone();
if predicate(&current) {
return current;
}
status.changed().await.expect("status channel");
}
})
.await
.expect("heartbeat status timeout")
}
async fn assert_credential_failure(config: HeartbeatConfig, server: &TestServer, expected: &str) {
let shutdown = CancellationToken::new();
let runtime = spawn_heartbeat_runtime(Some(config), &shutdown, summary)
.expect("start runtime")
.expect("configured runtime");
let mut status = runtime.status();
assert!(matches!(
wait_for(&mut status, |status| matches!(status, HeartbeatStatus::Failed { .. })).await,
HeartbeatStatus::Failed { reason } if reason.contains(expected)
));
assert!(server.seen.lock().expect("seen lock").is_empty());
runtime.shutdown().await;
}
#[tokio::test]
async fn connect_config_absent_starts_no_task() {
let shutdown = CancellationToken::new();
let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let sampled = calls.clone();
let runtime = spawn_heartbeat_runtime(None, &shutdown, move || {
sampled.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
summary()
})
.expect("absent config");
assert!(runtime.is_none());
tokio::task::yield_now().await;
assert_eq!(calls.load(std::sync::atomic::Ordering::Relaxed), 0);
}
#[tokio::test]
async fn duplicate_runtime_is_rejected_without_a_second_task() {
let pki = TestPki::new();
let mut reply = Reply::ok("2026-08-22T01:02:03Z");
reply.delay = Duration::from_secs(5);
let server = server(&pki, vec![reply]).await;
let temp = tempfile::tempdir().expect("tempdir");
let shutdown = CancellationToken::new();
let config = config(&temp, &pki, &server);
let runtime = spawn_heartbeat_runtime(Some(config.clone()), &shutdown, summary)
.expect("first runtime")
.expect("configured runtime");
assert!(matches!(
spawn_heartbeat_runtime(Some(config), &shutdown, summary),
Err(rustfs::connect::HeartbeatError::AlreadyRunning)
));
runtime.shutdown().await;
}
#[tokio::test(flavor = "current_thread")]
async fn dropped_runtime_keeps_the_lock_until_its_task_stops() {
let pki = TestPki::new();
let server = server(&pki, vec![Reply::ok("2026-08-22T01:02:03Z")]).await;
let temp = tempfile::tempdir().expect("tempdir");
let shutdown = CancellationToken::new();
let config = config(&temp, &pki, &server);
let runtime = spawn_heartbeat_runtime(Some(config.clone()), &shutdown, summary)
.expect("first runtime")
.expect("configured runtime");
drop(runtime);
assert!(matches!(
spawn_heartbeat_runtime(Some(config.clone()), &shutdown, summary),
Err(rustfs::connect::HeartbeatError::AlreadyRunning)
));
let replacement = tokio::time::timeout(Duration::from_secs(3), async {
loop {
match spawn_heartbeat_runtime(Some(config.clone()), &shutdown, summary) {
Ok(Some(runtime)) => break runtime,
Err(rustfs::connect::HeartbeatError::AlreadyRunning) => tokio::task::yield_now().await,
Ok(None) => panic!("configured replacement returned no runtime"),
Err(error) => panic!("unexpected replacement error: {error}"),
}
}
})
.await
.expect("dropped runtime releases its lock after stopping");
replacement.shutdown().await;
}
#[tokio::test]
async fn corrupt_persisted_state_is_rejected_before_network_delivery() {
let pki = TestPki::new();
let server = server(&pki, vec![Reply::ok("2026-08-22T01:02:03Z")]).await;
let temp = tempfile::tempdir().expect("tempdir");
let shutdown = CancellationToken::new();
let config = config(&temp, &pki, &server);
let directory = config.state_path.parent().expect("state directory");
fs::create_dir_all(directory).expect("create state directory");
fs::write(
&config.state_path,
br#"{"nextSequence":0,"pending":{"protocolVersion":"v1","requestId":"550e8400-e29b-41d4-a716-446655440000","agentVersion":"rustfs-agent/1.0.0-rc.3","capabilities":["heartbeat"],"sequence":0,"clientTime":"2026-08-22T01:02:03Z","coarseNodeSummary":{"total":0,"healthy":0,"degraded":0}}}"#,
)
.expect("write corrupt state");
private_mode(&config.state_path);
let runtime = spawn_heartbeat_runtime(Some(config), &shutdown, summary)
.expect("start runtime")
.expect("configured runtime");
let mut status = runtime.status();
assert!(matches!(
wait_for(&mut status, |status| matches!(status, HeartbeatStatus::Failed { .. })).await,
HeartbeatStatus::Failed { reason } if reason.contains("violates the protocol invariants")
));
assert!(server.seen.lock().expect("seen lock").is_empty());
runtime.shutdown().await;
}
#[tokio::test]
async fn invalid_stored_resource_name_is_rejected_before_network_delivery() {
let pki = TestPki::new();
let server = server(&pki, vec![]).await;
let temp = tempfile::tempdir().expect("tempdir");
let config = config(&temp, &pki, &server);
rewrite_credential(&temp, |credential| {
credential.name = format!("organizations/{ORGANIZATION_UID}/clusters/not-a-uuid/clusterDevices/{DEVICE_UID}");
});
assert_credential_failure(config, &server, "wrong device identity").await;
}
#[tokio::test]
async fn invalid_stored_protocol_is_rejected_before_network_delivery() {
let pki = TestPki::new();
let server = server(&pki, vec![]).await;
let temp = tempfile::tempdir().expect("tempdir");
let config = config(&temp, &pki, &server);
rewrite_credential(&temp, |credential| credential.protocol_version = "v2".to_owned());
assert_credential_failure(config, &server, "wrong device identity").await;
}
#[tokio::test]
async fn stored_certificate_key_mismatch_is_rejected_before_network_delivery() {
let pki = TestPki::new();
let server = server(&pki, vec![]).await;
let temp = tempfile::tempdir().expect("tempdir");
let now = OffsetDateTime::now_utc();
let (identity_store, credential_store) =
pki.stores_with_certificate(&temp, now - time::Duration::hours(1), now + time::Duration::hours(23), false);
let config = config_with_stores(&temp, &pki, &server, identity_store, credential_store);
assert_credential_failure(config, &server, "different device key").await;
}
#[tokio::test]
async fn expired_stored_certificate_is_rejected_before_network_delivery() {
let pki = TestPki::new();
let server = server(&pki, vec![]).await;
let temp = tempfile::tempdir().expect("tempdir");
let now = OffsetDateTime::now_utc();
let (identity_store, credential_store) =
pki.stores_with_certificate(&temp, now - time::Duration::days(2), now - time::Duration::days(1), true);
let config = config_with_stores(&temp, &pki, &server, identity_store, credential_store);
assert_credential_failure(config, &server, "not currently valid").await;
}
#[tokio::test]
async fn sends_only_l0_fields_and_accepts_additive_response_fields() {
let pki = TestPki::new();
let server = server(&pki, vec![Reply::ok("2038-01-19T03:14:07Z")]).await;
let temp = tempfile::tempdir().expect("tempdir");
let shutdown = CancellationToken::new();
let runtime = spawn_heartbeat_runtime(Some(config(&temp, &pki, &server)), &shutdown, summary)
.expect("start runtime")
.expect("configured runtime");
let mut status = runtime.status();
assert_eq!(
wait_for(&mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await,
HeartbeatStatus::Online {
server_time: "2038-01-19T03:14:07Z".to_owned()
}
);
runtime.shutdown().await;
let seen = server.seen.lock().expect("seen lock");
let request = &seen[0];
let mut keys = request
.as_object()
.expect("heartbeat object")
.keys()
.map(String::as_str)
.collect::<Vec<_>>();
keys.sort_unstable();
assert_eq!(
keys,
[
"agentVersion",
"capabilities",
"clientTime",
"coarseNodeSummary",
"protocolVersion",
"requestId",
"sequence"
]
);
assert_eq!(request["capabilities"], json!(["heartbeat"]));
assert_eq!(request["coarseNodeSummary"], json!({"total": 8, "healthy": 7, "degraded": 1}));
assert_ne!(request["clientTime"], "2038-01-19T03:14:07Z");
assert!(request.get("authorization").is_none());
}
#[tokio::test]
async fn restart_replays_pending_request_then_advances_sequence() {
let pki = TestPki::new();
let first_server = server(&pki, vec![Reply::error(StatusCode::SERVICE_UNAVAILABLE)]).await;
let temp = tempfile::tempdir().expect("tempdir");
let shutdown = CancellationToken::new();
let first_config = config(&temp, &pki, &first_server);
let runtime = spawn_heartbeat_runtime(Some(first_config.clone()), &shutdown, summary)
.expect("start runtime")
.expect("configured runtime");
let mut status = runtime.status();
wait_for(&mut status, |status| matches!(status, HeartbeatStatus::BackingOff { .. })).await;
runtime.shutdown().await;
let first = first_server.seen.lock().expect("seen lock")[0].clone();
drop(first_server);
let second_server = server(&pki, vec![Reply::ok("2026-08-22T01:02:03Z"), Reply::ok("2026-08-22T01:02:04Z")]).await;
let mut second_config = first_config;
second_config.endpoint = second_server.endpoint.clone();
let runtime = spawn_heartbeat_runtime(Some(second_config), &shutdown, summary)
.expect("restart runtime")
.expect("configured runtime");
tokio::time::timeout(Duration::from_secs(3), async {
while second_server.seen.lock().expect("seen lock").len() < 2 {
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("two heartbeats");
runtime.shutdown().await;
let seen = second_server.seen.lock().expect("seen lock");
assert_eq!(seen[0]["requestId"], first["requestId"]);
assert_eq!(seen[0]["sequence"], first["sequence"]);
assert_ne!(seen[1]["requestId"], seen[0]["requestId"]);
assert_eq!(seen[1]["sequence"].as_u64(), seen[0]["sequence"].as_u64().map(|value| value + 1));
}
#[tokio::test]
async fn retry_after_is_respected_with_the_local_upper_bound() {
let pki = TestPki::new();
let mut reply = Reply::error(StatusCode::TOO_MANY_REQUESTS);
reply.retry_after = Some("300");
let server = server(&pki, vec![reply]).await;
let temp = tempfile::tempdir().expect("tempdir");
let shutdown = CancellationToken::new();
let runtime = spawn_heartbeat_runtime(Some(config(&temp, &pki, &server)), &shutdown, summary)
.expect("start runtime")
.expect("configured runtime");
let mut status = runtime.status();
assert_eq!(
wait_for(&mut status, |status| matches!(status, HeartbeatStatus::BackingOff { .. })).await,
HeartbeatStatus::BackingOff {
delay: Duration::from_millis(80)
}
);
runtime.shutdown().await;
}
#[tokio::test]
async fn disconnects_use_exponential_backoff_with_a_cap() {
let pki = TestPki::new();
let server = server(
&pki,
vec![
Reply::error(StatusCode::SERVICE_UNAVAILABLE),
Reply::error(StatusCode::SERVICE_UNAVAILABLE),
Reply::error(StatusCode::SERVICE_UNAVAILABLE),
],
)
.await;
let temp = tempfile::tempdir().expect("tempdir");
let shutdown = CancellationToken::new();
let runtime = spawn_heartbeat_runtime(Some(config(&temp, &pki, &server)), &shutdown, summary)
.expect("start runtime")
.expect("configured runtime");
let mut status = runtime.status();
for delay in [20, 40, 80] {
assert_eq!(
wait_for(&mut status, |status| {
matches!(status, HeartbeatStatus::BackingOff { delay: observed } if *observed == Duration::from_millis(delay))
})
.await,
HeartbeatStatus::BackingOff {
delay: Duration::from_millis(delay)
}
);
}
runtime.shutdown().await;
}
#[tokio::test]
async fn revoked_credential_stops_and_exposes_local_status() {
let pki = TestPki::new();
let mut reply = Reply::error(StatusCode::UNAUTHORIZED);
reply.body = json!({"details": [{"reason": "CREDENTIAL_REVOKED"}]});
let server = server(&pki, vec![reply]).await;
let temp = tempfile::tempdir().expect("tempdir");
let shutdown = CancellationToken::new();
let runtime = spawn_heartbeat_runtime(Some(config(&temp, &pki, &server)), &shutdown, summary)
.expect("start runtime")
.expect("configured runtime");
let mut status = runtime.status();
assert_eq!(
wait_for(&mut status, |status| matches!(status, HeartbeatStatus::AuthenticationStopped { .. })).await,
HeartbeatStatus::AuthenticationStopped {
status: 401,
reason: Some("CREDENTIAL_REVOKED".to_owned())
}
);
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(server.seen.lock().expect("seen lock").len(), 1);
runtime.shutdown().await;
}
#[tokio::test]
async fn shutdown_cancels_an_in_flight_request() {
let pki = TestPki::new();
let mut reply = Reply::ok("2026-08-22T01:02:03Z");
reply.delay = Duration::from_secs(5);
let server = server(&pki, vec![reply]).await;
let temp = tempfile::tempdir().expect("tempdir");
let shutdown = CancellationToken::new();
let runtime = spawn_heartbeat_runtime(Some(config(&temp, &pki, &server)), &shutdown, summary)
.expect("start runtime")
.expect("configured runtime");
tokio::time::timeout(Duration::from_secs(3), async {
while server.seen.lock().expect("seen lock").is_empty() {
tokio::task::yield_now().await;
}
})
.await
.expect("request reached server");
tokio::time::timeout(Duration::from_millis(250), runtime.shutdown())
.await
.expect("cancellable shutdown");
}
#[test]
fn consumes_the_frozen_heartbeat_fixtures() {
let registry: Value =
serde_json::from_str(include_str!("../../protocol/agent/v1/fixtures/fixture-sets.json")).expect("fixture registry");
let heartbeat = registry["sets"]
.as_array()
.expect("fixture sets")
.iter()
.find(|set| set["name"] == "heartbeat")
.expect("heartbeat fixture set");
assert_eq!(heartbeat["status"], "populated");
let valid: Value =
serde_json::from_str(include_str!("../../protocol/agent/v1/fixtures/heartbeat/valid.json")).expect("valid fixture");
assert_eq!(valid["request"]["protocolVersion"], "v1");
let overflow: Value =
serde_json::from_str(include_str!("../../protocol/agent/v1/fixtures/heartbeat/overflow.json")).expect("overflow fixture");
assert_eq!(overflow["expected"]["httpStatus"], 422);
}
#[cfg(unix)]
fn private_mode(path: &Path) {
use std::os::unix::fs::PermissionsExt as _;
fs::set_permissions(path, fs::Permissions::from_mode(0o600)).expect("private mode");
}
#[cfg(not(unix))]
fn private_mode(_path: &Path) {}