fix(ilm): harden tier transition failure boundaries (#5031)

* fix(tier): fence generation-scoped operations

Refs rustfs/backlog#1354

Co-Authored-By: heihutu <heihutu@gmail.com>

* fix(ilm): verify transition upload streams

Refs rustfs/backlog#1353

Co-Authored-By: heihutu <heihutu@gmail.com>

* test(ecstore): expand transition fault matrix

Refs rustfs/backlog#1355

Co-Authored-By: heihutu <heihutu@gmail.com>

---------

Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
houseme
2026-07-19 18:48:32 +08:00
committed by GitHub
parent 83d73b34f3
commit 21049401fa
25 changed files with 5808 additions and 423 deletions
+4 -2
View File
@@ -67,7 +67,9 @@ pub mod bucket {
}
pub mod tier_delete_journal {
pub use crate::bucket::lifecycle::tier_delete_journal::persist_tier_delete_journal_entry;
pub use crate::bucket::lifecycle::tier_delete_journal::{
persist_tier_delete_journal_entry, record_tier_delete_journal_backend_identity,
};
}
pub mod tier_last_day_stats {
@@ -408,7 +410,7 @@ pub mod tier {
pub use crate::services::tier::tier::{
ERR_TIER_BACKEND_IN_USE, ERR_TIER_BACKEND_NOT_EMPTY, ERR_TIER_INVALID_CONFIG, ERR_TIER_MISSING_CREDENTIALS,
ERR_TIER_TYPE_UNSUPPORTED, TIER_CONFIG_FILE, TIER_CONFIG_FORMAT, TIER_CONFIG_V1, TIER_CONFIG_VERSION, TierConfigMgr,
is_err_config_not_found, try_migrate_tiering_config,
TierConfigUpdateError, is_err_config_not_found, try_migrate_tiering_config,
};
}
@@ -29,7 +29,7 @@ use crate::bucket::lifecycle::replication_sink::{
use crate::bucket::lifecycle::tier_delete_journal::{process_tier_delete_journal_entry, run_tier_delete_journal_recovery_loop};
use crate::bucket::lifecycle::tier_free_version_recovery::{DEFAULT_FREE_VERSION_RECOVERY_LIMIT, recover_tier_free_versions};
use crate::bucket::lifecycle::tier_last_day_stats::{DailyAllTierStats, LastDayTierStats};
use crate::bucket::lifecycle::tier_sweeper::{Jentry, delete_object_from_remote_tier_idempotent};
use crate::bucket::lifecycle::tier_sweeper::{Jentry, delete_object_from_remote_tier_idempotent_with_manager_and_identity};
use crate::bucket::versioning_sys::BucketVersioningSys;
use crate::client::object_api_utils::new_getobjectreader;
use crate::disk::error::DiskError;
@@ -38,7 +38,10 @@ use crate::error::Error;
use crate::error::StorageError;
use crate::error::{error_resp_to_object_err, is_err_object_not_found, is_err_version_not_found, is_network_or_host_down};
use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions};
use crate::services::tier::warm_backend::WarmBackendGetOpts;
use crate::services::tier::{
tier::{TierConfigMgr, tier_destination_id_from_metadata},
warm_backend::WarmBackendGetOpts,
};
use crate::set_disk::{MAX_PARTS_COUNT, RUSTFS_MULTIPART_BUCKET_KEY, RUSTFS_MULTIPART_OBJECT_KEY, SetDisks};
use crate::storage_api_contracts::{
lifecycle::ExpirationOptions,
@@ -402,6 +405,36 @@ impl ExpiryOp for FreeVersionTask {
}
}
async fn delete_free_version_remote_object(
oi: &ObjectInfo,
tier_config_mgr: &Arc<RwLock<TierConfigMgr>>,
) -> Result<(), std::io::Error> {
let identity = tier_destination_id_from_metadata(&oi.user_defined)?
.ok_or_else(|| std::io::Error::other("tier free-version has no durable backend identity"))?;
delete_object_from_remote_tier_idempotent_with_manager_and_identity(
&oi.transitioned_object.name,
&oi.transitioned_object.version_id,
&oi.transitioned_object.tier,
identity,
tier_config_mgr,
)
.await?;
Ok(())
}
async fn delete_free_version_remote_object_then<T, F, Fut>(
oi: &ObjectInfo,
tier_config_mgr: &Arc<RwLock<TierConfigMgr>>,
delete_local: F,
) -> Result<T, std::io::Error>
where
F: FnOnce() -> Fut,
Fut: std::future::Future<Output = T>,
{
delete_free_version_remote_object(oi, tier_config_mgr).await?;
Ok(delete_local().await)
}
struct NewerNoncurrentTask {
bucket: String,
versions: Vec<ObjectToDelete>,
@@ -680,13 +713,60 @@ impl ExpiryState {
else if v.as_any().is::<FreeVersionTask>() {
let v = v.as_any().downcast_ref::<FreeVersionTask>().expect("FreeVersionTask downcast failed");
let oi = v.0.clone();
if let Err(err) = delete_object_from_remote_tier_idempotent(
&oi.transitioned_object.name,
&oi.transitioned_object.version_id,
&oi.transitioned_object.tier,
)
.await
{
let cleanup = delete_free_version_remote_object_then(&oi, &api.tier_config_mgr(), || async {
let mut fi = FileInfo {
name: oi.name.clone(),
version_id: oi.version_id,
deleted: true,
..Default::default()
};
fi.set_tier_free_version();
let mut deleted_locally = false;
for pool in api.pools.iter() {
let set = pool.get_disks_by_key(&oi.name);
match set.delete_object_version(&oi.bucket, &oi.name, &fi, false).await {
Ok(()) => {
deleted_locally = true;
break;
}
Err(err) if is_err_version_not_found(&err) || is_err_object_not_found(&err) => continue,
Err(err) => {
debug!(
event = EVENT_LIFECYCLE_WORKER_STATE,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
bucket = %oi.bucket,
object = %oi.name,
remote_object = %oi.transitioned_object.name,
remote_version_id = %oi.transitioned_object.version_id,
tier = %oi.transitioned_object.tier,
error = ?err,
reason = "local_free_version_delete_failed",
"Lifecycle worker failed local free-version cleanup"
);
break;
}
}
}
if !deleted_locally {
debug!(
event = EVENT_LIFECYCLE_WORKER_STATE,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
bucket = %oi.bucket,
object = %oi.name,
remote_object = %oi.transitioned_object.name,
remote_version_id = %oi.transitioned_object.version_id,
tier = %oi.transitioned_object.tier,
reason = "local_free_version_missing",
"Lifecycle worker could not find transitioned free version locally"
);
}
})
.await;
if let Err(err) = cleanup {
debug!(
bucket = %oi.bucket,
object = %oi.name,
@@ -702,57 +782,6 @@ impl ExpiryState {
);
continue;
}
let mut fi = FileInfo {
name: oi.name.clone(),
version_id: oi.version_id,
deleted: true,
..Default::default()
};
fi.set_tier_free_version();
let mut deleted_locally = false;
for pool in api.pools.iter() {
let set = pool.get_disks_by_key(&oi.name);
match set.delete_object_version(&oi.bucket, &oi.name, &fi, false).await {
Ok(()) => {
deleted_locally = true;
break;
}
Err(err) if is_err_version_not_found(&err) || is_err_object_not_found(&err) => continue,
Err(err) => {
debug!(
event = EVENT_LIFECYCLE_WORKER_STATE,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
bucket = %oi.bucket,
object = %oi.name,
remote_object = %oi.transitioned_object.name,
remote_version_id = %oi.transitioned_object.version_id,
tier = %oi.transitioned_object.tier,
error = ?err,
reason = "local_free_version_delete_failed",
"Lifecycle worker failed local free-version cleanup"
);
break;
}
}
}
if !deleted_locally {
debug!(
event = EVENT_LIFECYCLE_WORKER_STATE,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
bucket = %oi.bucket,
object = %oi.name,
remote_object = %oi.transitioned_object.name,
remote_version_id = %oi.transitioned_object.version_id,
tier = %oi.transitioned_object.tier,
reason = "local_free_version_missing",
"Lifecycle worker could not find transitioned free version locally"
);
}
}
else {
//info!("Invalid work type - {:?}", v);
@@ -2444,8 +2473,27 @@ pub async fn get_transitioned_object_reader(
opts: &ObjectOptions,
) -> Result<GetObjectReader, std::io::Error> {
let tier_config_mgr = runtime_sources::tier_config_mgr_handle();
let mut tier_config_mgr = tier_config_mgr.write().await;
let tgt_client = match tier_config_mgr.get_driver(&oi.transitioned_object.tier).await {
get_transitioned_object_reader_with_tier_manager(bucket, object, rs, h, oi, opts, &tier_config_mgr).await
}
pub(crate) async fn get_transitioned_object_reader_with_tier_manager(
bucket: &str,
object: &str,
rs: &Option<HTTPRangeSpec>,
h: &HeaderMap,
oi: &ObjectInfo,
opts: &ObjectOptions,
tier_config_mgr: &Arc<RwLock<TierConfigMgr>>,
) -> Result<GetObjectReader, std::io::Error> {
let expected_identity = tier_destination_id_from_metadata(&oi.user_defined)?;
let lease = match expected_identity {
Some(identity) => {
TierConfigMgr::acquire_operation_lease_for_backend_identity(tier_config_mgr, &oi.transitioned_object.tier, identity)
.await
}
None => TierConfigMgr::acquire_operation_lease(tier_config_mgr, &oi.transitioned_object.tier).await,
};
let tgt_client = match lease {
Ok(d) => d,
Err(err) => return Err(std::io::Error::other(err)),
};
@@ -3081,6 +3129,8 @@ mod tests {
select_restore_s3_location, should_defer_date_expiry_for_recent_config_update,
should_reuse_lifecycle_delete_replication_state, transitioned_cleanup_tuple, transitioned_object_delete_opts,
};
#[cfg(feature = "test-util")]
use super::{delete_free_version_remote_object_then, get_transitioned_object_reader_with_tier_manager};
use crate::bucket::lifecycle::bucket_lifecycle_audit::LcEventSrc;
use crate::bucket::lifecycle::replication_sink::{
ReplicateDecision, ReplicateTargetDecision, ReplicationStatusType, VersionPurgeStatusType,
@@ -3121,6 +3171,233 @@ mod tests {
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
#[cfg(feature = "test-util")]
#[tokio::test]
async fn free_version_remote_delete_requires_persisted_destination_identity() {
let manager = crate::services::tier::tier::TierConfigMgr::new();
let old_backend = crate::services::tier::test_util::register_mock_tier(&manager, "WARM").await;
let mut conflicting_sys = HashMap::new();
rustfs_utils::http::metadata_compat::insert_bytes(
&mut conflicting_sys,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_STATUS,
crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.as_bytes().to_vec(),
);
conflicting_sys.insert(
format!(
"{}{}",
rustfs_utils::http::metadata_compat::RUSTFS_INTERNAL_PREFIX,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID
),
b"0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef".to_vec(),
);
conflicting_sys.insert(
format!(
"{}{}",
rustfs_utils::http::metadata_compat::MINIO_INTERNAL_PREFIX,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID
),
b"abcdef0123456789abcdef0123456789abcdef0123456789abcdef0123456789".to_vec(),
);
let conflicting_meta = rustfs_filemeta::MetaObject {
meta_sys: conflicting_sys,
..Default::default()
};
let mut free_version_info = rustfs_filemeta::FileInfo::new("object", 2, 2);
free_version_info.set_tier_free_version_id(&Uuid::new_v4().to_string());
assert_eq!(
conflicting_meta
.init_free_version(&free_version_info)
.expect_err("conflicting persisted identities must not create an executable free-version"),
rustfs_filemeta::Error::FileCorrupt
);
assert_eq!(old_backend.remove_count().await, 0);
let old_identity = crate::services::tier::tier::TierConfigMgr::acquire_operation_lease(&manager, "WARM")
.await
.expect("old tier lease should be available")
.backend_identity();
let mut oi = ObjectInfo::default();
oi.transitioned_object.tier = "WARM".to_string();
oi.transitioned_object.name = "remote/object".to_string();
oi.transitioned_object.version_id = "remote-version".to_string();
let local_delete_calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let legacy_err = delete_free_version_remote_object_then(&oi, &manager, {
let local_delete_calls = Arc::clone(&local_delete_calls);
move || async move {
local_delete_calls.fetch_add(1, Ordering::Relaxed);
}
})
.await
.expect_err("legacy free-version without identity must be retained");
assert!(legacy_err.to_string().contains("no durable backend identity"));
assert_eq!(local_delete_calls.load(Ordering::Relaxed), 0);
let mut invalid_metadata = HashMap::new();
rustfs_utils::http::metadata_compat::insert_str(
&mut invalid_metadata,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
"not-a-backend-identity".to_string(),
);
oi.user_defined = Arc::new(invalid_metadata);
let invalid_err = delete_free_version_remote_object_then(&oi, &manager, {
let local_delete_calls = Arc::clone(&local_delete_calls);
move || async move {
local_delete_calls.fetch_add(1, Ordering::Relaxed);
}
})
.await
.expect_err("free-version with an invalid identity must be retained");
assert!(invalid_err.to_string().contains("invalid length"));
assert_eq!(local_delete_calls.load(Ordering::Relaxed), 0);
let mut metadata = HashMap::new();
rustfs_utils::http::metadata_compat::insert_str(
&mut metadata,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
rustfs_utils::crypto::hex(old_identity),
);
oi.user_defined = Arc::new(metadata.clone());
delete_free_version_remote_object_then(&oi, &manager, {
let local_delete_calls = Arc::clone(&local_delete_calls);
move || async move {
local_delete_calls.fetch_add(1, Ordering::Relaxed);
}
})
.await
.expect("matching destination identity should allow idempotent remote cleanup");
assert_eq!(old_backend.remove_count().await, 1);
assert_eq!(local_delete_calls.load(Ordering::Relaxed), 1);
let mut single_prefix_metadata = HashMap::new();
single_prefix_metadata.insert(
format!(
"{}{}",
rustfs_utils::http::metadata_compat::MINIO_INTERNAL_PREFIX,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID
),
rustfs_utils::crypto::hex(old_identity),
);
oi.user_defined = Arc::new(single_prefix_metadata);
delete_free_version_remote_object_then(&oi, &manager, {
let local_delete_calls = Arc::clone(&local_delete_calls);
move || async move {
local_delete_calls.fetch_add(1, Ordering::Relaxed);
}
})
.await
.expect("single-prefix legacy identity should remain compatible");
assert_eq!(old_backend.remove_count().await, 2);
assert_eq!(local_delete_calls.load(Ordering::Relaxed), 2);
let new_backend = crate::services::tier::test_util::register_mock_tier(&manager, "WARM").await;
let new_identity = crate::services::tier::tier::TierConfigMgr::acquire_operation_lease(&manager, "WARM")
.await
.expect("rebound tier lease should be available")
.backend_identity();
let mut conflicting_metadata = HashMap::from([(
rustfs_utils::http::metadata_compat::internal_key_rustfs(
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
),
rustfs_utils::crypto::hex(new_identity),
)]);
conflicting_metadata.insert(
format!(
"{}{}",
rustfs_utils::http::metadata_compat::MINIO_INTERNAL_PREFIX,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID
),
rustfs_utils::crypto::hex(old_identity),
);
oi.user_defined = Arc::new(conflicting_metadata);
let conflict_err = delete_free_version_remote_object_then(&oi, &manager, {
let local_delete_calls = Arc::clone(&local_delete_calls);
move || async move {
local_delete_calls.fetch_add(1, Ordering::Relaxed);
}
})
.await
.expect_err("conflicting compatibility identities must retain the free-version");
assert!(conflict_err.to_string().contains("compatibility keys conflict"));
assert_eq!(new_backend.remove_count().await, 0);
assert_eq!(local_delete_calls.load(Ordering::Relaxed), 2);
oi.user_defined = Arc::new(metadata);
let rebound_err = delete_free_version_remote_object_then(&oi, &manager, {
let local_delete_calls = Arc::clone(&local_delete_calls);
move || async move {
local_delete_calls.fetch_add(1, Ordering::Relaxed);
}
})
.await
.expect_err("same-name tier rebind must retain the old free-version");
assert!(rebound_err.to_string().contains("identity no longer matches"));
assert_eq!(new_backend.remove_count().await, 0);
assert_eq!(local_delete_calls.load(Ordering::Relaxed), 2);
}
#[cfg(feature = "test-util")]
#[tokio::test]
async fn transitioned_get_rejects_same_name_rebind_before_remote_io() {
let manager = crate::services::tier::tier::TierConfigMgr::new();
crate::services::tier::test_util::register_mock_tier(&manager, "WARM").await;
let old_identity = crate::services::tier::tier::TierConfigMgr::acquire_operation_lease(&manager, "WARM")
.await
.expect("old tier lease should be available")
.backend_identity();
let new_backend = crate::services::tier::test_util::register_mock_tier(&manager, "WARM").await;
let mut metadata = HashMap::new();
rustfs_utils::http::metadata_compat::insert_str(
&mut metadata,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
rustfs_utils::crypto::hex(old_identity),
);
let mut oi = ObjectInfo {
user_defined: Arc::new(metadata),
..Default::default()
};
oi.transitioned_object.tier = "WARM".to_string();
oi.transitioned_object.name = "remote/object".to_string();
oi.transitioned_object.version_id = "remote-version".to_string();
let err = match get_transitioned_object_reader_with_tier_manager(
"bucket",
"object",
&None,
&http::HeaderMap::new(),
&oi,
&ObjectOptions::default(),
&manager,
)
.await
{
Ok(_) => panic!("identity-bound GET must reject a same-name tier rebind"),
Err(err) => err,
};
assert!(err.to_string().contains("identity no longer matches"));
assert_eq!(new_backend.get_count().await, 0);
oi.user_defined = Arc::new(HashMap::new());
let err = match get_transitioned_object_reader_with_tier_manager(
"bucket",
"object",
&None,
&http::HeaderMap::new(),
&oi,
&ObjectOptions::default(),
&manager,
)
.await
{
Ok(_) => panic!("missing remote legacy object should return an error"),
Err(err) => err,
};
assert!(!err.to_string().is_empty());
assert_eq!(new_backend.get_count().await, 1);
}
/// Pins the expiry-event routing for transitioned objects
/// (rustfs/backlog#1302): restore-expiry events must set
/// `transition.expire_restored` (strip-restored-copy semantics, never a
@@ -3192,6 +3469,7 @@ mod tests {
obj_name: "remote/object".to_string(),
version_id: "remote-version".to_string(),
tier_name: "WARM".to_string(),
backend_identity: Some([1; 32]),
};
let err = state
@@ -3279,6 +3557,7 @@ mod tests {
obj_name: "remote/object".to_string(),
version_id: "remote-version".to_string(),
tier_name: "WARM".to_string(),
backend_identity: Some([1; 32]),
};
state
@@ -20,10 +20,11 @@ use tokio_util::sync::CancellationToken;
use tracing::{debug, warn};
use crate::bucket::lifecycle::config_boundary;
use crate::bucket::lifecycle::tier_sweeper::{Jentry, delete_object_from_remote_tier_idempotent};
use crate::bucket::lifecycle::tier_sweeper::{Jentry, delete_object_from_remote_tier_idempotent_with_manager_and_identity};
use crate::disk::RUSTFS_META_BUCKET;
use crate::error::{Error, Result};
use crate::object_api::{GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader};
use crate::services::tier::tier::tier_destination_id_from_metadata;
use crate::storage_api_contracts::{
list::ListOperations as _,
object::{DeletedObject, ObjectIO, ObjectOperations, ObjectToDelete},
@@ -37,7 +38,7 @@ const LOG_SUBSYSTEM_LIFECYCLE: &str = "lifecycle";
const EVENT_LIFECYCLE_TIER_DELETE_JOURNAL: &str = "lifecycle_tier_delete_journal";
pub const DEFAULT_TIER_DELETE_JOURNAL_RECOVERY_LIMIT: usize = 1_000;
const TIER_DELETE_JOURNAL_VERSION: u8 = 1;
const TIER_DELETE_JOURNAL_VERSION: u8 = 2;
const TIER_DELETE_JOURNAL_PREFIX: &str = "ilm/tier-delete-journal/";
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
@@ -47,22 +48,26 @@ struct PersistedTierDeleteJournalEntry {
obj_name: String,
version_id: String,
tier_name: String,
#[serde(default)]
backend_identity: Option<[u8; 32]>,
}
impl PersistedTierDeleteJournalEntry {
fn from_jentry(je: &Jentry) -> Self {
Self {
version: TIER_DELETE_JOURNAL_VERSION,
version: if je.backend_identity.is_some() {
TIER_DELETE_JOURNAL_VERSION
} else {
1
},
obj_name: je.obj_name.clone(),
version_id: je.version_id.clone(),
tier_name: je.tier_name.clone(),
backend_identity: je.backend_identity,
}
}
fn into_jentry(self) -> Result<Jentry> {
if self.version != TIER_DELETE_JOURNAL_VERSION {
return Err(Error::other(format!("unsupported tier delete journal version {}", self.version)));
}
// Empty `version_id` is a legal sentinel for objects transitioned to an
// unversioned remote tier (see CLAUDE.md: a tier version of `None`/`""`
// means the tier bucket is unversioned, so the remote delete is issued
@@ -71,10 +76,19 @@ impl PersistedTierDeleteJournalEntry {
if self.obj_name.is_empty() || self.tier_name.is_empty() {
return Err(Error::other("tier delete journal entry is incomplete"));
}
let backend_identity = match self.version {
1 => None,
TIER_DELETE_JOURNAL_VERSION => Some(
self.backend_identity
.ok_or_else(|| Error::other("tier delete journal v2 entry is missing its backend identity"))?,
),
version => return Err(Error::other(format!("unsupported tier delete journal version {version}"))),
};
Ok(Jentry {
obj_name: self.obj_name,
version_id: self.version_id,
tier_name: self.tier_name,
backend_identity,
})
}
}
@@ -95,6 +109,10 @@ pub(crate) fn tier_delete_journal_object_name(je: &Jentry) -> String {
hasher.update(je.obj_name.as_bytes());
hasher.update([0]);
hasher.update(je.version_id.as_bytes());
if let Some(backend_identity) = je.backend_identity {
hasher.update([0]);
hasher.update(backend_identity);
}
format!(
"{TIER_DELETE_JOURNAL_PREFIX}{}.json",
rustfs_utils::crypto::hex(hasher.finalize().as_slice())
@@ -112,6 +130,16 @@ fn encode_tier_delete_journal_entry(je: &Jentry) -> Result<Vec<u8>> {
.map_err(|err| Error::other(format!("encode tier delete journal failed: {err}")))
}
pub fn record_tier_delete_journal_backend_identity(
je: &mut Jentry,
metadata: &std::collections::HashMap<String, String>,
) -> std::io::Result<()> {
if let Some(identity) = tier_destination_id_from_metadata(metadata)? {
je.backend_identity = Some(identity);
}
Ok(())
}
pub async fn persist_tier_delete_journal_entry<S>(api: Arc<S>, je: &Jentry) -> std::io::Result<()>
where
S: ObjectIO<
@@ -148,7 +176,17 @@ where
}
pub async fn process_tier_delete_journal_entry(api: Arc<ECStore>, je: &Jentry) -> std::io::Result<()> {
delete_object_from_remote_tier_idempotent(&je.obj_name, &je.version_id, &je.tier_name).await?;
let backend_identity = je
.backend_identity
.ok_or_else(|| std::io::Error::other("legacy tier delete journal has no durable backend identity"))?;
delete_object_from_remote_tier_idempotent_with_manager_and_identity(
&je.obj_name,
&je.version_id,
&je.tier_name,
backend_identity,
&api.tier_config_mgr(),
)
.await?;
remove_tier_delete_journal_entry(api, je).await
}
@@ -218,6 +256,21 @@ pub async fn recover_tier_delete_journal_entries(
}
};
if je.backend_identity.is_none() {
stats.failed += 1;
warn!(
event = EVENT_LIFECYCLE_TIER_DELETE_JOURNAL,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
journal_object = %object.name,
remote_object = %je.obj_name,
remote_version_id = %je.version_id,
tier = %je.tier_name,
"Legacy tier delete journal entry has no durable backend identity and will be retained"
);
continue;
}
match process_tier_delete_journal_entry(api.clone(), &je).await {
Ok(()) => stats.deleted += 1,
Err(err) => {
@@ -281,7 +334,10 @@ pub async fn run_tier_delete_journal_recovery_loop(api: Arc<ECStore>, cancel_tok
#[cfg(test)]
mod tests {
use super::{decode_tier_delete_journal_entry, encode_tier_delete_journal_entry, tier_delete_journal_object_name};
use super::{
decode_tier_delete_journal_entry, encode_tier_delete_journal_entry, record_tier_delete_journal_backend_identity,
tier_delete_journal_object_name,
};
use crate::bucket::lifecycle::tier_sweeper::Jentry;
fn journal_entry() -> Jentry {
@@ -289,6 +345,7 @@ mod tests {
obj_name: "remote/object".to_string(),
version_id: "remote-version".to_string(),
tier_name: "WARM".to_string(),
backend_identity: Some([7; 32]),
}
}
@@ -302,6 +359,7 @@ mod tests {
assert_eq!(decoded.obj_name, je.obj_name);
assert_eq!(decoded.version_id, je.version_id);
assert_eq!(decoded.tier_name, je.tier_name);
assert_eq!(decoded.backend_identity, je.backend_identity);
}
#[test]
@@ -317,6 +375,63 @@ mod tests {
assert!(!first.contains("remote/object"));
}
#[test]
fn tier_delete_journal_paths_separate_legacy_and_backend_identities() {
let mut legacy = journal_entry();
legacy.backend_identity = None;
let mut backend_a = journal_entry();
backend_a.backend_identity = Some([1; 32]);
let mut backend_b = journal_entry();
backend_b.backend_identity = Some([2; 32]);
assert_eq!(
tier_delete_journal_object_name(&legacy),
"ilm/tier-delete-journal/5ba6a7eb6338412b771613a6845a42ae5b8e26b5d201323eb01b38c5b42ff300.json"
);
assert_ne!(tier_delete_journal_object_name(&legacy), tier_delete_journal_object_name(&backend_a));
assert_ne!(tier_delete_journal_object_name(&backend_a), tier_delete_journal_object_name(&backend_b));
}
#[test]
fn tier_delete_journal_v2_requires_backend_identity() {
let payload = br#"{"version":2,"obj_name":"remote/object","version_id":"v1","tier_name":"WARM"}"#;
let err = decode_tier_delete_journal_entry(payload).expect_err("v2 entry without identity must fail closed");
assert!(err.to_string().contains("backend identity"));
}
#[test]
fn tier_delete_journal_uses_persisted_transition_destination_identity() {
let mut je = journal_entry();
je.backend_identity = None;
let identity = [9_u8; 32];
let mut metadata = std::collections::HashMap::new();
rustfs_utils::http::metadata_compat::insert_str(
&mut metadata,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
rustfs_utils::crypto::hex(identity),
);
record_tier_delete_journal_backend_identity(&mut je, &metadata).expect("persisted transition identity should decode");
let encoded = encode_tier_delete_journal_entry(&je).expect("identity-bound journal should encode");
let decoded = decode_tier_delete_journal_entry(&encoded).expect("identity-bound journal should decode");
assert_eq!(decoded.backend_identity, Some(identity));
}
#[test]
fn tier_delete_journal_without_transition_identity_stays_legacy() {
let mut je = journal_entry();
je.backend_identity = None;
let encoded = encode_tier_delete_journal_entry(&je).expect("legacy journal should remain encodable");
let persisted: serde_json::Value = serde_json::from_slice(&encoded).expect("journal JSON should decode");
assert_eq!(persisted["version"], 1);
assert!(persisted["backend_identity"].is_null());
}
#[test]
fn tier_delete_journal_rejects_incomplete_entry() {
let payload = br#"{"version":1,"obj_name":"","version_id":"v1","tier_name":"WARM"}"#;
@@ -339,6 +454,7 @@ mod tests {
assert_eq!(decoded.obj_name, "remote/object");
assert!(decoded.version_id.is_empty());
assert_eq!(decoded.tier_name, "WARM");
assert_eq!(decoded.backend_identity, None);
}
#[test]
@@ -23,6 +23,7 @@ use crate::bucket::lifecycle::bucket_lifecycle_ops::ExpiryOp;
use crate::bucket::lifecycle::lifecycle::{self, ObjectOpts};
use crate::bucket::lifecycle::tier_delete_journal::persist_tier_delete_journal_entry;
use crate::client::signer_error::error_chain_contains_signer_header_marker;
use crate::services::tier::tier::{TierConfigMgr, TierDestinationId, TierOperationLease};
use crate::storage_api_contracts::lifecycle::TransitionedObject;
use crate::store::ECStore;
use rustfs_utils::get_env_usize;
@@ -247,6 +248,7 @@ impl ObjSweeper {
obj_name: self.remote_object.clone(),
version_id: self.transition_version_id.clone(),
tier_name: self.transition_tier.clone(),
backend_identity: None,
});
}
None
@@ -281,6 +283,7 @@ pub struct Jentry {
pub(crate) obj_name: String,
pub(crate) version_id: String,
pub(crate) tier_name: String,
pub(crate) backend_identity: Option<TierDestinationId>,
}
impl ExpiryOp for Jentry {
@@ -312,6 +315,27 @@ async fn delete_object_from_remote_tier_raw(obj_name: &str, rv_id: &str, tier_na
return result;
}
let tier_config_mgr = runtime_sources::tier_config_mgr_handle();
delete_object_from_remote_tier_raw_with_manager(obj_name, rv_id, tier_name, &tier_config_mgr).await
}
async fn delete_object_from_remote_tier_raw_with_manager(
obj_name: &str,
rv_id: &str,
tier_name: &str,
tier_config_mgr: &Arc<tokio::sync::RwLock<TierConfigMgr>>,
) -> Result<(), std::io::Error> {
let lease = TierConfigMgr::acquire_operation_lease(&tier_config_mgr, tier_name)
.await
.map_err(std::io::Error::other)?;
delete_object_from_remote_tier_raw_with_lease(obj_name, rv_id, &lease).await
}
async fn delete_object_from_remote_tier_raw_with_lease(
obj_name: &str,
rv_id: &str,
lease: &TierOperationLease,
) -> Result<(), std::io::Error> {
if remote_delete_breaker_is_open(Instant::now()).await {
metrics::counter!(METRIC_DELETE_REMOTE_BREAKER_TOTAL).increment(1);
return Err(std::io::Error::other(ERR_REMOTE_DELETE_BREAKER_OPEN));
@@ -323,13 +347,7 @@ async fn delete_object_from_remote_tier_raw(obj_name: &str, rv_id: &str, tier_na
.map_err(|_| std::io::Error::other(ERR_REMOTE_DELETE_LIMITER_CLOSED))?;
let _inflight = RemoteDeleteInflightGuard::new();
let tier_config_mgr = runtime_sources::tier_config_mgr_handle();
let mut config_mgr = tier_config_mgr.write().await;
let w = match config_mgr.get_driver(tier_name).await {
Ok(w) => w,
Err(e) => return Err(std::io::Error::other(e)),
};
w.remove(obj_name, rv_id).await
lease.remove(obj_name, rv_id).await
}
#[cfg(test)]
@@ -364,6 +382,36 @@ pub async fn delete_object_from_remote_tier_idempotent(
}
}
pub(crate) async fn delete_object_from_remote_tier_idempotent_with_manager_and_identity(
obj_name: &str,
rv_id: &str,
tier_name: &str,
backend_identity: TierDestinationId,
tier_config_mgr: &Arc<tokio::sync::RwLock<TierConfigMgr>>,
) -> Result<RemoteTierDeleteOutcome, std::io::Error> {
let lease = TierConfigMgr::acquire_operation_lease_for_backend_identity(tier_config_mgr, tier_name, backend_identity)
.await
.map_err(std::io::Error::other)?;
delete_object_from_remote_tier_with_lease_idempotent(obj_name, rv_id, &lease).await
}
pub(crate) async fn delete_object_from_remote_tier_with_lease_idempotent(
obj_name: &str,
rv_id: &str,
lease: &TierOperationLease,
) -> Result<RemoteTierDeleteOutcome, std::io::Error> {
match delete_object_from_remote_tier_raw_with_lease(obj_name, rv_id, lease).await {
Ok(()) => Ok(RemoteTierDeleteOutcome::Deleted),
Err(err) if is_remote_tier_not_found_error(&err) => Ok(RemoteTierDeleteOutcome::AlreadyRemoved),
Err(err) => {
if should_record_remote_delete_failure(&err) {
record_remote_delete_failure(&err, Instant::now()).await;
}
Err(err)
}
}
}
pub(crate) fn is_remote_tier_not_found_error(err: &std::io::Error) -> bool {
let message = err.to_string();
message.contains("NoSuchKey")
@@ -401,6 +449,7 @@ pub fn transitioned_force_delete_journal_entry(transitioned: &TransitionedObject
obj_name: transitioned.name.clone(),
version_id: transitioned.version_id.clone(),
tier_name: transitioned.tier.clone(),
backend_identity: None,
})
}
@@ -410,7 +459,8 @@ mod test {
use super::{
ERR_REMOTE_DELETE_BREAKER_OPEN, ERR_REMOTE_DELETE_LIMITER_CLOSED, REMOTE_TIER_DELETE_TEST_HOOK, RemoteDeleteBreaker,
RemoteTierDeleteOutcome, delete_object_from_remote_tier_idempotent, is_remote_tier_not_found_error,
RemoteTierDeleteOutcome, delete_object_from_remote_tier_idempotent,
delete_object_from_remote_tier_idempotent_with_manager_and_identity, is_remote_tier_not_found_error,
is_signer_header_error, should_record_remote_delete_failure,
};
use std::io::{Error, ErrorKind};
@@ -506,6 +556,31 @@ mod test {
assert!(err.to_string().contains("driver not found"));
}
#[cfg(feature = "test-util")]
#[tokio::test]
async fn journal_delete_rejects_backend_identity_mismatch() {
let manager = crate::services::tier::tier::TierConfigMgr::new();
crate::services::tier::test_util::register_mock_tier(&manager, "WARM").await;
let lease = crate::services::tier::tier::TierConfigMgr::acquire_operation_lease(&manager, "WARM")
.await
.expect("test tier lease should be available");
let mut mismatched = lease.backend_identity();
mismatched[0] ^= 1;
drop(lease);
let err = delete_object_from_remote_tier_idempotent_with_manager_and_identity(
"remote/object",
"remote-version",
"WARM",
mismatched,
&manager,
)
.await
.expect_err("journal recovery must fail closed when the tier name was rebound");
assert!(err.to_string().contains("identity no longer matches"));
}
#[test]
fn breaker_opens_at_threshold_and_recovers_after_window() {
let mut breaker = RemoteDeleteBreaker::new(3, Duration::from_secs(30));
+6 -1
View File
@@ -552,7 +552,12 @@ pub(crate) async fn initialize_local_disk_maps(
}
pub(crate) async fn init_tier_config_mgr(store: Arc<ECStore>) -> Result<()> {
get_global_tier_config_mgr().write().await.init(store).await
let handle = get_global_tier_config_mgr();
TierConfigMgr::reload_handle(&handle, store.clone()).await?;
if setup_is_dist_erasure().await {
tokio::spawn(TierConfigMgr::refresh_tier_config_handle(handle, store));
}
Ok(())
}
#[cfg(test)]
+34 -5
View File
@@ -140,6 +140,8 @@ pub struct MockStoredObject {
struct MockWarmBackendInner {
objects: Mutex<HashMap<String, MockStoredObject>>,
faults: Mutex<FaultConfig>,
put_read_limit: Mutex<Option<usize>>,
put_remote_version: Mutex<Option<String>>,
op_log: Mutex<Vec<MockWarmOp>>,
put_versions: Mutex<Vec<(String, String)>>,
remove_versions: Mutex<Vec<(String, String)>>,
@@ -232,6 +234,18 @@ impl MockWarmBackend {
*self.inner.faults.lock().await = FaultConfig::default();
}
/// Limit how many body bytes a successful mock PUT consumes. `None` drains
/// the complete body. This models a backend that incorrectly accepts a
/// truncated stream while still returning success.
pub async fn set_put_read_limit(&self, limit: Option<usize>) {
*self.inner.put_read_limit.lock().await = limit;
}
/// Override the remote version returned by subsequent successful PUTs.
pub async fn set_put_remote_version(&self, remote_version: Option<String>) {
*self.inner.put_remote_version.lock().await = remote_version;
}
async fn precondition(&self) -> Result<(), std::io::Error> {
let (latency, error) = {
let faults = self.inner.faults.lock().await;
@@ -379,7 +393,13 @@ impl MockWarmBackend {
// ---- internal helpers -----------------------------------------------
async fn put_bytes(&self, object: &str, bytes: Vec<u8>, metadata: HashMap<String, String>) -> String {
let remote_version_id = Uuid::new_v4().to_string();
let remote_version_id = self
.inner
.put_remote_version
.lock()
.await
.clone()
.unwrap_or_else(|| Uuid::new_v4().to_string());
self.inner.objects.lock().await.insert(
object.to_string(),
MockStoredObject {
@@ -392,11 +412,18 @@ impl MockWarmBackend {
}
async fn read_bytes(&self, reader: ReaderImpl) -> Result<Vec<u8>, std::io::Error> {
let limit = *self.inner.put_read_limit.lock().await;
match reader {
ReaderImpl::Body(bytes) => Ok(bytes.to_vec()),
ReaderImpl::Body(bytes) => Ok(bytes.slice(..limit.unwrap_or(bytes.len()).min(bytes.len())).to_vec()),
ReaderImpl::ObjectBody(mut reader) => {
let mut buf = Vec::new();
reader.stream.read_to_end(&mut buf).await?;
if let Some(limit) = limit {
let limit =
u64::try_from(limit).map_err(|_| std::io::Error::other("mock PUT read limit exceeds u64::MAX"))?;
reader.stream.take(limit).read_to_end(&mut buf).await?;
} else {
reader.stream.read_to_end(&mut buf).await?;
}
Ok(buf)
}
}
@@ -519,7 +546,7 @@ impl WarmBackend for MockWarmBackend {
async fn in_use(&self) -> Result<bool, std::io::Error> {
self.precondition().await?;
self.record(MockWarmOp::InUse).await;
Ok(false)
Ok(!self.inner.objects.lock().await.is_empty())
}
}
@@ -558,7 +585,9 @@ pub async fn register_mock_tier_backend(handle: &Arc<RwLock<TierConfigMgr>>, tie
..Default::default()
},
);
tier_config_mgr.driver_cache.insert(tier_name.to_string(), Box::new(backend));
tier_config_mgr
.install_test_driver(tier_name, Box::new(backend))
.expect("mock tier driver should install");
}
/// The transition-state tuple read from an on-disk `xl.meta`, plus the object's
File diff suppressed because it is too large Load Diff
@@ -238,6 +238,23 @@ impl Clone for TierConfig {
#[allow(dead_code)]
impl TierConfig {
pub(crate) fn clone_with_credentials(&self) -> Self {
Self {
version: self.version.clone(),
tier_type: self.tier_type.clone(),
name: self.name.clone(),
s3: self.s3.clone(),
aliyun: self.aliyun.clone(),
tencent: self.tencent.clone(),
huaweicloud: self.huaweicloud.clone(),
azure: self.azure.clone(),
gcs: self.gcs.clone(),
r2: self.r2.clone(),
rustfs: self.rustfs.clone(),
minio: self.minio.clone(),
}
}
fn endpoint(&self) -> String {
match self.tier_type {
TierType::S3 => self.s3.as_ref().map(|s| s.endpoint.clone()).unwrap_or_default(),
@@ -65,7 +65,15 @@ pub struct WarmBackendGetOpts {
#[async_trait::async_trait]
pub trait WarmBackend {
/// Return `Ok` only after the backend has consumed the complete declared
/// body and its storage service has acknowledged the PUT. The built-in S3
/// family uses the transition client's declared-length request plus
/// Content-MD5 for multipart parts, while GCS materializes the body before
/// awaiting its buffered write response. Test backends may deliberately
/// violate this contract to exercise transition compensation.
async fn put(&self, object: &str, r: ReaderImpl, length: i64) -> Result<String, std::io::Error>;
/// The same completion contract as [`WarmBackend::put`] applies when
/// metadata is attached.
async fn put_with_meta(
&self,
object: &str,
+3 -1
View File
@@ -97,7 +97,7 @@ use crate::storage_api_contracts::{
use crate::store::utils::is_reserved_or_invalid_bucket;
use crate::{
bucket::lifecycle::bucket_lifecycle_ops::{
LifecycleOps, gen_transition_objname, get_transitioned_object_reader, put_restore_opts,
LifecycleOps, gen_transition_objname, get_transitioned_object_reader_with_tier_manager, put_restore_opts,
},
cache_value::metacache_set::{ListPathRawOptions, list_path_raw},
config::storageclass,
@@ -628,6 +628,8 @@ pub(crate) use ops::object::body_cache_plaintext_len;
mod read;
mod replication;
pub(crate) mod shard_source;
#[cfg(all(test, feature = "test-util"))]
mod transition_matrix_tests;
pub use ops::heal_walk::HealWalkVersion;
@@ -1607,6 +1607,74 @@ mod tests {
.await
}
#[tokio::test]
#[serial(metadata_cache_invalidation_probe)]
async fn complete_multipart_generation_retires_cached_snapshot() {
use crate::storage_api_contracts::multipart::MultipartOperations as _;
use crate::storage_api_contracts::object::ObjectIO as _;
let (_temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
let bucket = "multipart-metadata-generation-bucket";
let object = "object";
for disk in &disk_stores {
disk.make_volume(bucket).await.expect("bucket volume should be created");
}
let mut initial_reader = PutObjReader::from_vec(b"old multipart body".to_vec());
set_disks
.put_object(bucket, object, &mut initial_reader, &ObjectOptions::default())
.await
.expect("initial object should be written");
set_disks
.get_object_fileinfo(bucket, object, &ObjectOptions::default(), true, false)
.await
.expect("initial metadata should resolve");
let generation = set_disks
.get_object_metadata_cache_generation(bucket, object)
.expect("metadata cache generation should be active");
let retired_key = GetObjectMetadataCacheKey::new(bucket, object, generation);
assert!(set_disks.get_object_metadata_cache.get(&retired_key).await.is_some());
let upload = set_disks
.new_multipart_upload(bucket, object, &ObjectOptions::default())
.await
.expect("multipart upload should be created");
let payload = vec![9u8; 4096];
let payload_len = i64::try_from(payload.len()).expect("test payload length should fit i64");
let mut part_reader = PutObjReader::new(
HashReader::from_stream(Cursor::new(payload), payload_len, payload_len, None, None, false)
.expect("part hash reader should be created"),
);
let part = set_disks
.put_object_part(bucket, object, &upload.upload_id, 1, &mut part_reader, &ObjectOptions::default())
.await
.expect("multipart part should be written");
let invalidations = MetadataCacheInvalidationProbe::install(bucket, object);
set_disks
.clone()
.complete_multipart_upload(
bucket,
object,
&upload.upload_id,
vec![CompletePart {
part_num: part.part_num,
etag: part.etag,
..Default::default()
}],
&ObjectOptions::default(),
)
.await
.expect("multipart completion should succeed");
assert_eq!(
invalidations.count(),
2,
"multipart completion must invalidate before mutation and after commit"
);
set_disks.get_object_metadata_cache.run_pending_tasks().await;
assert!(set_disks.get_object_metadata_cache.get(&retired_key).await.is_none());
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn multipart_upload_read_lock_waits_for_upload_writer() {
File diff suppressed because it is too large Load Diff
+139 -5
View File
@@ -2540,20 +2540,27 @@ mod metadata_cache_tests {
}
#[tokio::test]
async fn get_object_metadata_cache_invalidation_removes_object_entry() {
async fn metadata_cache_per_key_invalidation_physically_reclaims_retired_generation() {
let set = new_metadata_cache_test_set().await;
let fi = valid_test_fileinfo("object");
let generation = set.get_object_metadata_cache_generation("bucket", "object");
set.cache_get_object_fileinfo(("bucket", "object"), generation, &fi, std::slice::from_ref(&fi), &[], 0)
let generation = set
.get_object_metadata_cache_generation("bucket", "object")
.expect("metadata cache generation should be active");
let retired_key = GetObjectMetadataCacheKey::new("bucket", "object", generation);
set.cache_get_object_fileinfo(("bucket", "object"), Some(generation), &fi, std::slice::from_ref(&fi), &[], 0)
.await;
set.get_object_metadata_cache.run_pending_tasks().await;
assert!(set.cached_get_object_fileinfo("bucket", "object").await.is_some());
assert_eq!(set.get_object_metadata_cache.entry_count(), 1);
set.invalidate_get_object_metadata_cache("bucket", "object").await;
set.get_object_metadata_cache.run_pending_tasks().await;
assert!(
set.cached_get_object_fileinfo("bucket", "object").await.is_none(),
"explicit invalidation must remove the cached object metadata"
set.get_object_metadata_cache.get(&retired_key).await.is_none(),
"per-key invalidation must physically remove the retired generation"
);
assert_eq!(set.get_object_metadata_cache.entry_count(), 0);
}
#[tokio::test]
@@ -2704,6 +2711,77 @@ mod metadata_cache_tests {
);
}
#[tokio::test]
async fn metadata_cache_generation_isolated_between_set_instances() {
let first = new_metadata_cache_test_set().await;
let second = new_metadata_cache_test_set().await;
let fi = valid_test_fileinfo("object");
let first_generation = first.get_object_metadata_cache_generation("bucket", "object");
let second_generation = second.get_object_metadata_cache_generation("bucket", "object");
first
.cache_get_object_fileinfo(("bucket", "object"), first_generation, &fi, std::slice::from_ref(&fi), &[], 0)
.await;
second
.cache_get_object_fileinfo(("bucket", "object"), second_generation, &fi, std::slice::from_ref(&fi), &[], 0)
.await;
first.invalidate_get_object_metadata_cache("bucket", "object").await;
assert!(first.cached_get_object_fileinfo("bucket", "object").await.is_none());
assert!(second.cached_get_object_fileinfo("bucket", "object").await.is_some());
assert_eq!(second.get_object_metadata_cache_generation("bucket", "object"), second_generation);
}
#[tokio::test]
async fn metadata_cache_cached_hash_collision_preserves_full_identity() {
let set = new_metadata_cache_test_set().await;
let generation = set
.get_object_metadata_cache_generation("bucket-a", "object-a")
.expect("metadata cache generation should be active");
let first_key = GetObjectMetadataCacheKey::new("bucket-a", "object-a", generation);
let second_key = GetObjectMetadataCacheKey {
bucket: Arc::from("bucket-b"),
object: Arc::from("object-b"),
generation: generation.value,
hash: generation.hash,
};
let first_fi = valid_test_fileinfo("object-a");
let second_fi = valid_test_fileinfo("object-b");
let entry = |fi: FileInfo| {
Arc::new(GetObjectMetadataCacheEntry {
created_at: Instant::now(),
parts_metadata: vec![fi.clone()],
fi,
online_disks: Vec::new(),
read_quorum: 0,
})
};
set.get_object_metadata_cache.insert(first_key.clone(), entry(first_fi)).await;
set.get_object_metadata_cache
.insert(second_key.clone(), entry(second_fi))
.await;
assert_eq!(
set.get_object_metadata_cache
.get(&first_key)
.await
.expect("first colliding entry should remain addressable")
.fi
.name,
"object-a"
);
assert_eq!(
set.get_object_metadata_cache
.get(&second_key)
.await
.expect("second colliding entry should remain addressable")
.fi
.name,
"object-b"
);
}
#[tokio::test]
async fn metadata_cache_generation_overflow_fails_closed() {
let set = new_metadata_cache_test_set().await;
@@ -2743,6 +2821,62 @@ mod metadata_cache_tests {
assert!(set.cached_get_object_fileinfo("bucket", "object").await.is_none());
}
#[tokio::test]
async fn metadata_cache_invalidate_all_physically_reclaims_retired_generations() {
let set = new_metadata_cache_test_set().await;
let mut retired_keys = Vec::new();
for object in ["object-a", "object-b", "object-c"] {
let fi = valid_test_fileinfo(object);
let generation = set
.get_object_metadata_cache_generation("bucket", object)
.expect("metadata cache generation should be active");
retired_keys.push(GetObjectMetadataCacheKey::new("bucket", object, generation));
set.cache_get_object_fileinfo(("bucket", object), Some(generation), &fi, std::slice::from_ref(&fi), &[], 0)
.await;
}
set.get_object_metadata_cache.run_pending_tasks().await;
assert_eq!(set.get_object_metadata_cache.entry_count(), 3);
set.invalidate_all_get_object_metadata_cache();
set.get_object_metadata_cache.run_pending_tasks().await;
for key in retired_keys {
assert!(
set.get_object_metadata_cache.get(&key).await.is_none(),
"invalidate-all must physically remove every retired generation"
);
}
assert_eq!(set.get_object_metadata_cache.entry_count(), 0);
}
#[tokio::test]
async fn metadata_cache_invalidate_all_at_max_fails_closed_and_clears_entries() {
let set = new_metadata_cache_test_set().await;
let fi = valid_test_fileinfo("object");
let generation = set
.get_object_metadata_cache_generation("bucket", "object")
.expect("metadata cache generation should be active");
let retired_key = GetObjectMetadataCacheKey::new("bucket", "object", generation);
set.cache_get_object_fileinfo(("bucket", "object"), Some(generation), &fi, std::slice::from_ref(&fi), &[], 0)
.await;
for fence in set.get_object_metadata_cache_generations.iter() {
fence.store(u64::MAX, Ordering::Release);
}
set.invalidate_all_get_object_metadata_cache();
set.get_object_metadata_cache.run_pending_tasks().await;
assert!(
set.get_object_metadata_cache_generations
.iter()
.all(|fence| fence.load(Ordering::Acquire) == u64::MAX),
"invalidate-all must not wrap a saturated fence"
);
assert_eq!(set.get_object_metadata_cache_generation("bucket", "object"), None);
assert!(set.get_object_metadata_cache.get(&retired_key).await.is_none());
assert_eq!(set.get_object_metadata_cache.entry_count(), 0);
}
#[tokio::test]
async fn get_object_metadata_cache_prunes_when_capacity_is_reached() {
// moka handles capacity eviction automatically via the configured max_capacity.
@@ -0,0 +1,104 @@
// 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 super::*;
use crate::bucket::lifecycle::lifecycle::{TRANSITION_PENDING, TransitionOptions};
use crate::ecstore_validation_blackbox::make_local_set_disks;
use crate::services::tier::test_util::register_mock_tier;
use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations as _};
use tokio::io::AsyncReadExt;
async fn prime_metadata_generation(set_disks: &SetDisks, bucket: &str, object: &str) -> GetObjectMetadataCacheKey {
set_disks
.get_object_fileinfo(bucket, object, &ObjectOptions::default(), true, false)
.await
.expect("object metadata should resolve");
let generation = set_disks
.get_object_metadata_cache_generation(bucket, object)
.expect("metadata generation should be active");
let key = GetObjectMetadataCacheKey::new(bucket, object, generation);
assert!(
set_disks.get_object_metadata_cache.get(&key).await.is_some(),
"metadata read should publish the generation under test"
);
key
}
async fn assert_generation_reclaimed(set_disks: &SetDisks, key: &GetObjectMetadataCacheKey) {
set_disks.get_object_metadata_cache.run_pending_tasks().await;
assert!(
set_disks.get_object_metadata_cache.get(key).await.is_none(),
"metadata mutation must physically reclaim the prior generation"
);
}
#[tokio::test]
#[serial_test::serial]
async fn transition_and_restore_reclaim_prior_metadata_generations() {
let (_dirs, set_disks) = make_local_set_disks(4, 2).await;
let bucket = "transition-restore-generation-bucket";
let object = "object.bin";
let payload = vec![0x5au8; 1024 * 1024];
set_disks
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created");
let mut reader = PutObjReader::from_vec(payload.clone());
let original = set_disks
.put_object(bucket, object, &mut reader, &ObjectOptions::default())
.await
.expect("source object should be written");
let source_generation = prime_metadata_generation(&set_disks, bucket, object).await;
let tier_name = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase();
let backend = register_mock_tier(&set_disks.instance_ctx().tier_config_mgr(), &tier_name).await;
let transition_opts = ObjectOptions {
transition: TransitionOptions {
status: TRANSITION_PENDING.to_string(),
tier: tier_name,
etag: original.etag.clone().expect("source ETag should be present"),
..Default::default()
},
version_id: original.version_id.map(|version| version.to_string()),
mod_time: original.mod_time,
..Default::default()
};
set_disks
.transition_object(bucket, object, &transition_opts)
.await
.expect("transition should succeed");
assert_generation_reclaimed(&set_disks, &source_generation).await;
let transitioned_generation = prime_metadata_generation(&set_disks, bucket, object).await;
let mut restore_opts = ObjectOptions::default();
restore_opts.transition.restore_request.days = Some(1);
Arc::clone(&set_disks)
.restore_transitioned_object(bucket, object, &restore_opts)
.await
.expect("restore should succeed");
assert_generation_reclaimed(&set_disks, &transitioned_generation).await;
assert_eq!(backend.get_count().await, 1, "restore should read the remote candidate exactly once");
let mut restored = Vec::new();
set_disks
.get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default())
.await
.expect("restored object should be readable")
.stream
.read_to_end(&mut restored)
.await
.expect("restored body should drain");
assert_eq!(restored, payload);
assert_eq!(backend.get_count().await, 1, "restored GET should use the local copy");
}
+1 -1
View File
@@ -264,7 +264,7 @@ impl ECStore {
/// Get the tier config manager
pub fn tier_config_mgr(&self) -> Arc<tokio::sync::RwLock<crate::services::tier::tier::TierConfigMgr>> {
runtime_sources::global_tier_config_mgr()
self.ctx.tier_config_mgr()
}
/// Get the server configuration
+121 -2
View File
@@ -30,8 +30,9 @@ use crate::ChecksumInfo;
use rustfs_utils::HashAlgorithm;
use rustfs_utils::http::{
SUFFIX_CRC, SUFFIX_FREE_VERSION, SUFFIX_INLINE_DATA, SUFFIX_PURGESTATUS, SUFFIX_TIER_FV_ID, SUFFIX_TIER_FV_MARKER,
SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID,
contains_key_bytes, get_bytes, has_internal_suffix, insert_bytes, is_internal_key, remove_bytes, strip_internal_prefix,
SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITION_TIER_DESTINATION_ID, SUFFIX_TRANSITIONED_OBJECTNAME,
SUFFIX_TRANSITIONED_VERSION_ID, contains_key_bytes, get_bytes, get_consistent_bytes, get_str, has_internal_suffix,
insert_bytes, is_internal_key, remove_bytes, strip_internal_prefix,
};
const MSGPACK_EXT8: u8 = 0xc7;
@@ -2403,6 +2404,9 @@ impl MetaObject {
);
}
insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITION_TIER, fi.transition_tier.as_bytes().to_vec());
if let Some(destination_id) = get_str(&fi.metadata, SUFFIX_TRANSITION_TIER_DESTINATION_ID) {
insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID, destination_id.into_bytes());
}
}
pub fn remove_restore_hdrs(&mut self) {
@@ -2467,6 +2471,16 @@ impl MetaObject {
insert_bytes(&mut delete_marker.meta_sys, suffix, v);
}
}
if contains_key_bytes(&self.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID) {
let destination_id = get_consistent_bytes(&self.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID)
.filter(|value| value.len() == 64 && value.iter().all(u8::is_ascii_hexdigit))
.ok_or(Error::FileCorrupt)?;
insert_bytes(
&mut delete_marker.meta_sys,
SUFFIX_TRANSITION_TIER_DESTINATION_ID,
destination_id.to_vec(),
);
}
return Ok((free_entry, true));
}
Ok((FileMetaVersion::default(), false))
@@ -4204,6 +4218,111 @@ mod tests {
assert!(matches!(err, Error::UuidParse(_)));
}
#[test]
fn meta_object_init_free_version_preserves_transition_destination_identity() {
let mut sys = HashMap::new();
insert_bytes(&mut sys, SUFFIX_TRANSITION_STATUS, TRANSITION_COMPLETE.as_bytes().to_vec());
let identity = b"0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef".to_vec();
insert_bytes(&mut sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID, identity.clone());
let obj = make_meta_object_with_sys(sys);
let mut fi = FileInfo::new("object", 2, 2);
fi.set_tier_free_version_id(&Uuid::new_v4().to_string());
let (free_version, created) = obj
.init_free_version(&fi)
.expect("free-version initialization should succeed");
let meta_sys = &free_version
.delete_marker
.expect("free-version should be a delete marker")
.meta_sys;
assert!(created);
assert_eq!(get_bytes(meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID), Some(identity));
assert!(meta_sys.contains_key(&format!(
"{}{}",
rustfs_utils::http::metadata_compat::RUSTFS_INTERNAL_PREFIX,
SUFFIX_TRANSITION_TIER_DESTINATION_ID
)));
assert!(meta_sys.contains_key(&format!(
"{}{}",
rustfs_utils::http::metadata_compat::MINIO_INTERNAL_PREFIX,
SUFFIX_TRANSITION_TIER_DESTINATION_ID
)));
}
#[test]
fn meta_object_init_free_version_accepts_single_prefix_transition_destination_identity() {
let mut sys = HashMap::new();
insert_bytes(&mut sys, SUFFIX_TRANSITION_STATUS, TRANSITION_COMPLETE.as_bytes().to_vec());
let identity = b"0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef".to_vec();
sys.insert(
format!(
"{}{}",
rustfs_utils::http::metadata_compat::MINIO_INTERNAL_PREFIX,
SUFFIX_TRANSITION_TIER_DESTINATION_ID
),
identity.clone(),
);
let obj = make_meta_object_with_sys(sys);
let mut fi = FileInfo::new("object", 2, 2);
fi.set_tier_free_version_id(&Uuid::new_v4().to_string());
let (free_version, created) = obj
.init_free_version(&fi)
.expect("single-prefix legacy destination identity should remain compatible");
let meta_sys = &free_version
.delete_marker
.expect("free-version should be a delete marker")
.meta_sys;
assert!(created);
assert_eq!(
get_consistent_bytes(meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID),
Some(identity.as_slice())
);
}
#[test]
fn meta_object_init_free_version_rejects_conflicting_or_invalid_transition_destination_identity() {
let valid_identity = b"0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef".to_vec();
for (rustfs_identity, minio_identity) in [
(
valid_identity,
b"abcdef0123456789abcdef0123456789abcdef0123456789abcdef0123456789".to_vec(),
),
(b"not-hex".to_vec(), b"not-hex".to_vec()),
(Vec::new(), Vec::new()),
] {
let mut sys = HashMap::new();
insert_bytes(&mut sys, SUFFIX_TRANSITION_STATUS, TRANSITION_COMPLETE.as_bytes().to_vec());
sys.insert(
format!(
"{}{}",
rustfs_utils::http::metadata_compat::RUSTFS_INTERNAL_PREFIX,
SUFFIX_TRANSITION_TIER_DESTINATION_ID
),
rustfs_identity,
);
sys.insert(
format!(
"{}{}",
rustfs_utils::http::metadata_compat::MINIO_INTERNAL_PREFIX,
SUFFIX_TRANSITION_TIER_DESTINATION_ID
),
minio_identity,
);
let obj = make_meta_object_with_sys(sys);
let mut fi = FileInfo::new("object", 2, 2);
fi.set_tier_free_version_id(&Uuid::new_v4().to_string());
assert_eq!(
obj.init_free_version(&fi)
.expect_err("unsafe destination identity must fail closed"),
Error::FileCorrupt
);
}
}
#[test]
fn delete_marker_decode_skips_unknown_fields_for_forward_compat() {
// A newer writer emits the three known fields plus an extra one. Decoding
+65
View File
@@ -38,6 +38,7 @@ pub const SUFFIX_TRANSITION_STATUS: &str = "transition-status";
pub const SUFFIX_TRANSITIONED_OBJECTNAME: &str = "transitioned-object";
pub const SUFFIX_TRANSITIONED_VERSION_ID: &str = "transitioned-versionID";
pub const SUFFIX_TRANSITION_TIER: &str = "transition-tier";
pub const SUFFIX_TRANSITION_TIER_DESTINATION_ID: &str = "transition-tier-destination-id";
pub const SUFFIX_FREE_VERSION: &str = "free-version";
pub const SUFFIX_PURGESTATUS: &str = "purgestatus";
pub const SUFFIX_REPLICA_STATUS: &str = "replica-status";
@@ -173,6 +174,28 @@ pub fn get_str(map: &HashMap<String, String>, suffix: &str) -> Option<String> {
.map(|(_, value)| value.clone())
}
fn get_consistent_value<'a, V: AsRef<[u8]>>(map: &'a HashMap<String, V>, suffix: &str) -> Option<&'a V> {
let (rustfs_key, minio_key) = both_keys(suffix);
let mut value = None;
for (key, candidate) in map {
if !key.eq_ignore_ascii_case(&rustfs_key) && !key.eq_ignore_ascii_case(&minio_key) {
continue;
}
if candidate.as_ref().is_empty() || value.is_some_and(|current: &V| current.as_ref() != candidate.as_ref()) {
return None;
}
value = Some(candidate);
}
value
}
/// Returns a non-empty value when every compatibility key present for `suffix` agrees.
/// A single RustFS or MinIO key is accepted for backward compatibility; conflicting or empty
/// values return `None` so callers at destructive boundaries can fail closed.
pub fn get_consistent_str<'a>(map: &'a HashMap<String, String>, suffix: &str) -> Option<&'a str> {
get_consistent_value(map, suffix).map(String::as_str)
}
pub fn contains_key_str(map: &HashMap<String, String>, suffix: &str) -> bool {
if with_internal_key(RUSTFS_INTERNAL_PREFIX, suffix, |k1| map.contains_key(k1)) {
return true;
@@ -206,6 +229,11 @@ pub fn get_bytes(map: &HashMap<String, Vec<u8>>, suffix: &str) -> Option<Vec<u8>
.or_else(|| with_internal_key(MINIO_INTERNAL_PREFIX, suffix, |k2| map.get(k2).cloned()))
}
/// Byte-valued counterpart of [`get_consistent_str`].
pub fn get_consistent_bytes<'a>(map: &'a HashMap<String, Vec<u8>>, suffix: &str) -> Option<&'a [u8]> {
get_consistent_value(map, suffix).map(Vec::as_slice)
}
pub fn contains_key_bytes(map: &HashMap<String, Vec<u8>>, suffix: &str) -> bool {
with_internal_key(RUSTFS_INTERNAL_PREFIX, suffix, |k1| map.contains_key(k1))
|| with_internal_key(MINIO_INTERNAL_PREFIX, suffix, |k2| map.contains_key(k2))
@@ -271,6 +299,43 @@ mod tests {
assert_eq!(get_str(&metadata, SUFFIX_TRANSITION_TIER).as_deref(), Some("rustfs-tier"));
}
#[test]
fn test_consistent_str_accepts_single_or_matching_values_and_rejects_conflicts() {
let rustfs_key = internal_key_rustfs(SUFFIX_TRANSITION_TIER_DESTINATION_ID);
let minio_key = format!("{MINIO_INTERNAL_PREFIX}{SUFFIX_TRANSITION_TIER_DESTINATION_ID}");
let mut metadata = HashMap::from([(rustfs_key, "identity-a".to_string())]);
assert_eq!(get_consistent_str(&metadata, SUFFIX_TRANSITION_TIER_DESTINATION_ID), Some("identity-a"));
metadata.insert(minio_key, "identity-a".to_string());
assert_eq!(get_consistent_str(&metadata, SUFFIX_TRANSITION_TIER_DESTINATION_ID), Some("identity-a"));
metadata.insert(
format!("{MINIO_INTERNAL_PREFIX}{SUFFIX_TRANSITION_TIER_DESTINATION_ID}"),
"identity-b".to_string(),
);
assert_eq!(get_consistent_str(&metadata, SUFFIX_TRANSITION_TIER_DESTINATION_ID), None);
}
#[test]
fn test_consistent_bytes_accepts_single_or_matching_values_and_rejects_conflicts() {
let rustfs_key = internal_key_rustfs(SUFFIX_TRANSITION_TIER_DESTINATION_ID);
let minio_key = format!("{MINIO_INTERNAL_PREFIX}{SUFFIX_TRANSITION_TIER_DESTINATION_ID}");
let mut metadata = HashMap::from([(rustfs_key, b"identity-a".to_vec())]);
assert_eq!(
get_consistent_bytes(&metadata, SUFFIX_TRANSITION_TIER_DESTINATION_ID),
Some(b"identity-a".as_slice())
);
metadata.insert(minio_key.clone(), b"identity-a".to_vec());
assert_eq!(
get_consistent_bytes(&metadata, SUFFIX_TRANSITION_TIER_DESTINATION_ID),
Some(b"identity-a".as_slice())
);
metadata.insert(minio_key, b"identity-b".to_vec());
assert_eq!(get_consistent_bytes(&metadata, SUFFIX_TRANSITION_TIER_DESTINATION_ID), None);
}
#[test]
fn test_bytes_lookup_falls_back_to_minio_key() {
let mut meta_sys =
+210 -189
View File
@@ -14,10 +14,11 @@
#![allow(unused_variables, unused_mut, unused_must_use)]
use crate::admin::runtime_sources::object_store_from_extensions;
use crate::admin::storage_api::runtime_sources::TierConfigMgr;
use crate::admin::storage_api::tier::{
AdminError, DailyAllTierStats, ERR_TIER_ALREADY_EXISTS, ERR_TIER_BACKEND_IN_USE, ERR_TIER_BACKEND_NOT_EMPTY,
ERR_TIER_CONNECT_ERR, ERR_TIER_INVALID_CREDENTIALS, ERR_TIER_MISSING_CREDENTIALS, ERR_TIER_NAME_NOT_UPPERCASE,
ERR_TIER_NOT_FOUND, ERR_TIER_RESERVED_NAME, TierConfig, TierCreds, TierType,
ERR_TIER_NOT_FOUND, ERR_TIER_RESERVED_NAME, TierConfig, TierConfigUpdateError, TierCreds, TierType,
};
use crate::{
admin::runtime_sources::{current_daily_tier_stats, current_notification_system, current_tier_config_handle},
@@ -29,8 +30,7 @@ use crate::{
server::{ADMIN_PREFIX, RemoteAddr},
storage::request_context::spawn_traced,
};
use http::Uri;
use http::{HeaderMap, StatusCode};
use http::{HeaderMap, StatusCode, Uri};
use hyper::Method;
use matchit::Params;
use percent_encoding::percent_decode_str;
@@ -100,6 +100,52 @@ fn spawn_transition_tier_config_propagation(action: &'static str) {
}
}
fn tier_mutation_error(
update_error: TierConfigUpdateError,
action: &'static str,
failure_code: &'static str,
) -> Result<AdminError, S3Error> {
match update_error {
TierConfigUpdateError::Load(err) => {
warn!(
event = EVENT_ADMIN_TIER_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_TIER,
action,
result = "reload_failed",
error = ?err,
"admin tier state"
);
Err(S3Error::with_message(
S3ErrorCode::Custom(failure_code.into()),
format!("tier reload failed. {err}"),
))
}
TierConfigUpdateError::Save(err) => {
warn!(
event = EVENT_ADMIN_TIER_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_TIER,
action,
result = "save_failed",
error = ?err,
"admin tier state"
);
Err(S3Error::with_message(S3ErrorCode::Custom(failure_code.into()), "tier save failed"))
}
TierConfigUpdateError::Mutation(err) | TierConfigUpdateError::Publish(err) => Ok(err),
}
}
fn tier_backend_in_use_response(err: &AdminError) -> Option<S3Error> {
(err.code == ERR_TIER_BACKEND_IN_USE.code)
.then(|| S3Error::with_message(S3ErrorCode::Custom("TierNameBackendInUse".into()), "tier backend is not empty"))
}
fn clear_tier_error_response(err: &AdminError) -> S3Error {
S3Error::with_message(S3ErrorCode::Custom("TierClearFailed".into()), format!("tier clear failed. {err}"))
}
fn resolve_tier_name(uri: &Uri, params: &Params<'_, '_>) -> S3Result<String> {
if let Some(tier) = params.get("tier") {
let decoded = percent_decode_str(tier)
@@ -315,87 +361,58 @@ impl Operation for AddTier {
return Err(s3_error!(InternalError, "object store is not initialized"));
};
{
let tier_config_mgr_handle = current_tier_config_handle();
let mut tier_config_mgr = tier_config_mgr_handle.write().await;
if let Err(err) = tier_config_mgr.reload(store).await {
let tier_config_mgr_handle = current_tier_config_handle();
if let Err(update_err) = TierConfigMgr::add_and_save(&tier_config_mgr_handle, store, args, force).await {
let err = tier_mutation_error(update_err, "add_tier", "TierAddFailed")?;
return if err.code == ERR_TIER_RESERVED_NAME.code {
warn!(
event = EVENT_ADMIN_TIER_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_TIER,
action = "add_tier",
result = "reload_failed",
tier_name = %tier_name_for_log,
result = "reserved_name_rejected",
"admin tier state"
);
Err(s3_error!(InvalidRequest, "Cannot use reserved tier name"))
} else if err.code == ERR_TIER_ALREADY_EXISTS.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierNameAlreadyExist".into()),
"tier name already exists",
))
} else if err.code == ERR_TIER_NAME_NOT_UPPERCASE.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierNameNotUppercase".into()),
"tier name must be uppercase",
))
} else if err.code == ERR_TIER_BACKEND_IN_USE.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierNameBackendInUse!".into()),
"tier backend is already in use",
))
} else if err.code == ERR_TIER_CONNECT_ERR.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierConnectError".into()),
"tier connectivity check failed",
))
} else if err.code == ERR_TIER_INVALID_CREDENTIALS.code {
Err(S3Error::with_message(S3ErrorCode::Custom(err.code.clone().into()), err.message))
} else {
warn!(
event = EVENT_ADMIN_TIER_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_TIER,
action = "add_tier",
tier_name = %tier_name_for_log,
result = "add_failed",
error = ?err,
"admin tier state"
);
return Err(S3Error::with_message(
Err(S3Error::with_message(
S3ErrorCode::Custom("TierAddFailed".into()),
format!("tier reload failed. {err}"),
));
}
if let Err(err) = tier_config_mgr.add(args, force).await {
return if err.code == ERR_TIER_RESERVED_NAME.code {
warn!(
event = EVENT_ADMIN_TIER_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_TIER,
action = "add_tier",
tier_name = %tier_name_for_log,
result = "reserved_name_rejected",
"admin tier state"
);
Err(s3_error!(InvalidRequest, "Cannot use reserved tier name"))
} else if err.code == ERR_TIER_ALREADY_EXISTS.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierNameAlreadyExist".into()),
"tier name already exists",
))
} else if err.code == ERR_TIER_NAME_NOT_UPPERCASE.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierNameNotUppercase".into()),
"tier name must be uppercase",
))
} else if err.code == ERR_TIER_BACKEND_IN_USE.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierNameBackendInUse!".into()),
"tier backend is already in use",
))
} else if err.code == ERR_TIER_CONNECT_ERR.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierConnectError".into()),
"tier connectivity check failed",
))
} else if err.code == ERR_TIER_INVALID_CREDENTIALS.code {
Err(S3Error::with_message(S3ErrorCode::Custom(err.code.clone().into()), err.message))
} else {
warn!(
event = EVENT_ADMIN_TIER_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_TIER,
action = "add_tier",
tier_name = %tier_name_for_log,
result = "add_failed",
error = ?err,
"admin tier state"
);
Err(S3Error::with_message(
S3ErrorCode::Custom("TierAddFailed".into()),
format!("tier add failed. {err}"),
))
};
}
if let Err(e) = tier_config_mgr.save().await {
warn!(
event = EVENT_ADMIN_TIER_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_TIER,
action = "add_tier",
result = "save_failed",
error = ?e,
"admin tier state"
);
return Err(S3Error::with_message(S3ErrorCode::Custom("TierAddFailed".into()), "tier save failed"));
}
format!("tier add failed. {err}"),
))
};
}
spawn_transition_tier_config_propagation("add");
@@ -473,61 +490,32 @@ impl Operation for EditTier {
return Err(s3_error!(InternalError, "object store is not initialized"));
};
{
let tier_config_mgr_handle = current_tier_config_handle();
let mut tier_config_mgr = tier_config_mgr_handle.write().await;
if let Err(err) = tier_config_mgr.reload(store).await {
let tier_config_mgr_handle = current_tier_config_handle();
if let Err(update_err) = TierConfigMgr::edit_and_save(&tier_config_mgr_handle, store, &tier_name, creds).await {
let err = tier_mutation_error(update_err, "edit_tier", "TierEditFailed")?;
return if err.code == ERR_TIER_NOT_FOUND.code {
Err(S3Error::with_message(S3ErrorCode::Custom("TierNotFound".into()), "tier not found"))
} else if err.code == ERR_TIER_MISSING_CREDENTIALS.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierMissingCredentials".into()),
"tier credentials are required",
))
} else {
warn!(
event = EVENT_ADMIN_TIER_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_TIER,
action = "edit_tier",
result = "reload_failed",
tier_name = %tier_name,
result = "edit_failed",
error = ?err,
"admin tier state"
);
return Err(S3Error::with_message(
Err(S3Error::with_message(
S3ErrorCode::Custom("TierEditFailed".into()),
format!("tier reload failed. {err}"),
));
}
if let Err(err) = tier_config_mgr.edit(&tier_name, creds).await {
return if err.code == ERR_TIER_NOT_FOUND.code {
Err(S3Error::with_message(S3ErrorCode::Custom("TierNotFound".into()), "tier not found"))
} else if err.code == ERR_TIER_MISSING_CREDENTIALS.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierMissingCredentials".into()),
"tier credentials are required",
))
} else {
warn!(
event = EVENT_ADMIN_TIER_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_TIER,
action = "edit_tier",
tier_name = %tier_name,
result = "edit_failed",
error = ?err,
"admin tier state"
);
Err(S3Error::with_message(
S3ErrorCode::Custom("TierEditFailed".into()),
format!("tier edit failed. {err}"),
))
};
}
if let Err(e) = tier_config_mgr.save().await {
warn!(
event = EVENT_ADMIN_TIER_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_TIER,
action = "edit_tier",
result = "save_failed",
error = ?e,
"admin tier state"
);
return Err(S3Error::with_message(S3ErrorCode::Custom("TierEditFailed".into()), "tier save failed"));
}
format!("tier edit failed. {err}"),
))
};
}
spawn_transition_tier_config_propagation("edit");
@@ -643,62 +631,34 @@ impl Operation for RemoveTier {
return Err(s3_error!(InternalError, "object store is not initialized"));
};
{
let tier_config_mgr_handle = current_tier_config_handle();
let mut tier_config_mgr = tier_config_mgr_handle.write().await;
if let Err(err) = tier_config_mgr.reload(store).await {
let tier_config_mgr_handle = current_tier_config_handle();
if let Err(update_err) = TierConfigMgr::remove_and_save(&tier_config_mgr_handle, store, &tier_name, force).await {
let err = tier_mutation_error(update_err, "remove_tier", "TierRemoveFailed")?;
return if err.code == ERR_TIER_NOT_FOUND.code {
Err(S3Error::with_message(S3ErrorCode::Custom("TierNotFound".into()), "tier not found"))
} else if err.code == ERR_TIER_BACKEND_NOT_EMPTY.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierNameBackendInUse".into()),
"tier backend is not empty",
))
} else if let Some(response) = tier_backend_in_use_response(&err) {
Err(response)
} else {
warn!(
event = EVENT_ADMIN_TIER_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_TIER,
action = "remove_tier",
result = "reload_failed",
tier_name = %tier_name,
result = "remove_failed",
error = ?err,
"admin tier state"
);
return Err(S3Error::with_message(
Err(S3Error::with_message(
S3ErrorCode::Custom("TierRemoveFailed".into()),
format!("tier reload failed. {err}"),
));
}
if let Err(err) = tier_config_mgr.remove(&tier_name, force).await {
return if err.code == ERR_TIER_NOT_FOUND.code {
Err(S3Error::with_message(S3ErrorCode::Custom("TierNotFound".into()), "tier not found"))
} else if err.code == ERR_TIER_BACKEND_NOT_EMPTY.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierNameBackendInUse".into()),
"tier backend is not empty",
))
} else {
warn!(
event = EVENT_ADMIN_TIER_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_TIER,
action = "remove_tier",
tier_name = %tier_name,
result = "remove_failed",
error = ?err,
"admin tier state"
);
Err(S3Error::with_message(
S3ErrorCode::Custom("TierRemoveFailed".into()),
format!("tier remove failed. {err}"),
))
};
}
if let Err(e) = tier_config_mgr.save().await {
warn!(
event = EVENT_ADMIN_TIER_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_TIER,
action = "remove_tier",
result = "save_failed",
error = ?e,
"admin tier state"
);
return Err(S3Error::with_message(S3ErrorCode::Custom("TierRemoveFailed".into()), "tier save failed"));
}
format!("tier remove failed. {err}"),
))
};
}
spawn_transition_tier_config_propagation("remove");
@@ -733,8 +693,9 @@ impl Operation for VerifyTier {
let tier = resolve_tier_name(&req.uri, &params)?;
let tier_config_mgr_handle = current_tier_config_handle();
let mut tier_config_mgr = tier_config_mgr_handle.write().await;
tier_config_mgr.verify(&tier).await.map_err(map_tier_verify_error)?;
TierConfigMgr::verify_without_manager_lock(&tier_config_mgr_handle, &tier)
.await
.map_err(map_tier_verify_error)?;
let mut header = HeaderMap::new();
header.insert(CONTENT_TYPE, "application/json".parse().expect("valid header value"));
@@ -921,10 +882,42 @@ impl Operation for ClearTier {
return Err(s3_error!(InvalidRequest, "invalid clear-tier confirmation token"));
};
let Some(store) = object_store_from_extensions(&req.extensions) else {
return Err(s3_error!(InternalError, "object store is not initialized"));
};
let tier_config_mgr_handle = current_tier_config_handle();
let mut tier_config_mgr = tier_config_mgr_handle.write().await;
//tier_config_mgr.reload(api);
if let Err(err) = tier_config_mgr.clear_tier(force).await {
if let Err(update_err) = TierConfigMgr::clear_and_save(&tier_config_mgr_handle, store, force).await {
let err = match update_err {
TierConfigUpdateError::Load(err) => {
warn!(
event = EVENT_ADMIN_TIER_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_TIER,
action = "clear_tier",
result = "reload_failed",
error = ?err,
"admin tier state"
);
return Err(S3Error::with_message(
S3ErrorCode::Custom("TierClearFailed".into()),
format!("tier clear failed. {err}"),
));
}
TierConfigUpdateError::Save(err) => {
warn!(
event = EVENT_ADMIN_TIER_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_TIER,
action = "clear_tier",
result = "save_failed",
error = ?err,
"admin tier state"
);
return Err(S3Error::with_message(S3ErrorCode::Custom("TierEditFailed".into()), "tier save failed"));
}
TierConfigUpdateError::Mutation(err) | TierConfigUpdateError::Publish(err) => err,
};
warn!(
event = EVENT_ADMIN_TIER_STATE,
component = LOG_COMPONENT_ADMIN,
@@ -934,22 +927,7 @@ impl Operation for ClearTier {
error = ?err,
"admin tier state"
);
return Err(S3Error::with_message(
S3ErrorCode::Custom("TierClearFailed".into()),
format!("tier clear failed. {err}"),
));
}
if let Err(e) = tier_config_mgr.save().await {
warn!(
event = EVENT_ADMIN_TIER_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_TIER,
action = "clear_tier",
result = "save_failed",
error = ?e,
"admin tier state"
);
return Err(S3Error::with_message(S3ErrorCode::Custom("TierEditFailed".into()), "tier save failed"));
return Err(clear_tier_error_response(&err));
}
let mut header = HeaderMap::new();
@@ -1106,6 +1084,49 @@ mod tests {
assert_eq!(mapped.message(), Some("tier verification failed. backend unavailable"));
}
#[test]
fn tier_mutation_error_preserves_reload_and_save_responses() {
let reload = tier_mutation_error(
TierConfigUpdateError::Load(std::io::Error::other("read failed")),
"add_tier",
"TierAddFailed",
)
.expect_err("reload error should map to an S3 response");
assert_eq!(reload.code(), &S3ErrorCode::Custom("TierAddFailed".into()));
assert_eq!(reload.message(), Some("tier reload failed. read failed"));
let save = tier_mutation_error(
TierConfigUpdateError::Save(std::io::Error::other("conditional write failed")),
"edit_tier",
"TierEditFailed",
)
.expect_err("save error should map to an S3 response");
assert_eq!(save.code(), &S3ErrorCode::Custom("TierEditFailed".into()));
assert_eq!(save.message(), Some("tier save failed"));
}
#[test]
fn clear_preserves_legacy_backend_not_empty_error_code() {
let response = clear_tier_error_response(&ERR_TIER_BACKEND_NOT_EMPTY);
assert_eq!(response.code(), &S3ErrorCode::Custom("TierClearFailed".into()));
assert!(
response
.message()
.is_some_and(|message| message.starts_with("tier clear failed."))
);
let response = clear_tier_error_response(&ERR_TIER_BACKEND_IN_USE);
assert_eq!(response.code(), &S3ErrorCode::Custom("TierClearFailed".into()));
assert!(
response
.message()
.is_some_and(|message| message.starts_with("tier clear failed."))
);
assert!(tier_backend_in_use_response(&ERR_TIER_BACKEND_NOT_EMPTY).is_none());
assert!(tier_backend_in_use_response(&ERR_TIER_NOT_FOUND).is_none());
}
#[test]
fn parse_clear_tier_query_rejects_unknown_duplicate_and_invalid_force() {
for raw in [
+2 -1
View File
@@ -110,6 +110,7 @@ pub(crate) type Result<T> = core::result::Result<T, Error>;
pub(crate) type TierConfig = ecstore_tier::tier_config::TierConfig;
pub(crate) type TierCreds = ecstore_tier::tier_admin::TierCreds;
pub(crate) type TierType = ecstore_tier::tier_config::TierType;
pub(crate) type TierConfigUpdateError = crate::storage::storage_api::TierConfigUpdateError;
pub(crate) mod runtime_sources {
pub(crate) type DailyAllTierStats = super::DailyAllTierStats;
@@ -575,6 +576,6 @@ pub(crate) mod tier {
pub(crate) use super::{
AdminError, DailyAllTierStats, ERR_TIER_ALREADY_EXISTS, ERR_TIER_BACKEND_IN_USE, ERR_TIER_BACKEND_NOT_EMPTY,
ERR_TIER_CONNECT_ERR, ERR_TIER_INVALID_CREDENTIALS, ERR_TIER_MISSING_CREDENTIALS, ERR_TIER_NAME_NOT_UPPERCASE,
ERR_TIER_NOT_FOUND, ERR_TIER_RESERVED_NAME, TierConfig, TierCreds, TierType,
ERR_TIER_NOT_FOUND, ERR_TIER_RESERVED_NAME, TierConfig, TierConfigUpdateError, TierCreds, TierType,
};
}
@@ -866,9 +866,6 @@ async fn get_transitioned_object_uses_remote_codec_fallback_path() {
.collect();
create_test_bucket(&ecstore, bucket.as_str()).await;
set_bucket_lifecycle_transition_with_tier(bucket.as_str(), &tier_name)
.await
.expect("Failed to set lifecycle configuration");
let uploaded = upload_test_object(&ecstore, bucket.as_str(), object, &payload).await;
let transition_opts = ObjectOptions {
+1 -1
View File
@@ -33,7 +33,7 @@ mod data_usage_snapshot_gating_test;
#[cfg(test)]
mod delete_objects_stat_gating_test;
#[cfg(test)]
mod gating_test_env;
pub(crate) mod gating_test_env;
#[cfg(test)]
mod lifecycle_transition_api_test;
#[cfg(test)]
+96 -1
View File
@@ -757,10 +757,11 @@ async fn enqueue_transitioned_delete_cleanup(
&existing.transitioned_object,
)
};
let Some(je) = je else {
let Some(mut je) = je else {
return Ok(());
};
tier_delete_journal::record_tier_delete_journal_backend_identity(&mut je, &existing.user_defined)?;
tier_delete_journal::persist_tier_delete_journal_entry(store, &je).await?;
let expiry_state = current_expiry_state_handle();
@@ -9250,6 +9251,100 @@ mod tests {
(store, context)
}
#[tokio::test]
#[serial_test::serial]
async fn transitioned_delete_cleanup_persists_identity_bound_and_legacy_journals() {
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 identity = [11_u8; 32];
let mut metadata = HashMap::new();
rustfs_utils::http::metadata_compat::insert_str(
&mut metadata,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
rustfs_utils::crypto::hex(identity),
);
let mut current = ObjectInfo {
user_defined: Arc::new(metadata),
..Default::default()
};
current.transitioned_object.status = lifecycle::TRANSITION_COMPLETE.to_string();
current.transitioned_object.tier = "WARM".to_string();
current.transitioned_object.name = "remote/identity-bound".to_string();
current.transitioned_object.version_id = "remote-version".to_string();
let journal_name = |remote_object: &str, backend_identity: Option<[u8; 32]>| {
use sha2::{Digest, Sha256};
let mut hasher = Sha256::new();
hasher.update(b"WARM");
hasher.update([0]);
hasher.update(remote_object.as_bytes());
hasher.update([0]);
hasher.update(b"remote-version");
if let Some(backend_identity) = backend_identity {
hasher.update([0]);
hasher.update(backend_identity);
}
format!("ilm/tier-delete-journal/{}.json", rustfs_utils::crypto::hex(hasher.finalize().as_slice()))
};
enqueue_transitioned_delete_cleanup(store.clone(), "bucket", "identity-bound", &ObjectOptions::default(), Some(&current))
.await
.expect("normal transitioned delete should persist an identity-bound journal");
let mut identity_bound = store
.get_object_reader(
".rustfs.sys",
&journal_name("remote/identity-bound", Some(identity)),
None,
http::HeaderMap::new(),
&ObjectOptions::default(),
)
.await
.expect("identity-bound journal should be readable");
let mut identity_bound_data = Vec::new();
tokio::io::AsyncReadExt::read_to_end(&mut identity_bound.stream, &mut identity_bound_data)
.await
.expect("identity-bound journal body should be readable");
let identity_bound: serde_json::Value =
serde_json::from_slice(&identity_bound_data).expect("identity-bound journal should decode as JSON");
assert_eq!(identity_bound["version"], serde_json::json!(2));
assert_eq!(identity_bound["backend_identity"], serde_json::json!(identity));
current.user_defined = Arc::new(HashMap::new());
current.transitioned_object.name = "remote/legacy".to_string();
enqueue_transitioned_delete_cleanup(
store.clone(),
"bucket",
"legacy",
&ObjectOptions {
delete_prefix: true,
..Default::default()
},
Some(&current),
)
.await
.expect("legacy force-delete cleanup should persist a fail-closed v1 journal");
let mut legacy = store
.get_object_reader(
".rustfs.sys",
&journal_name("remote/legacy", None),
None,
http::HeaderMap::new(),
&ObjectOptions::default(),
)
.await
.expect("legacy journal should be readable");
let mut legacy_data = Vec::new();
tokio::io::AsyncReadExt::read_to_end(&mut legacy.stream, &mut legacy_data)
.await
.expect("legacy journal body should be readable");
let legacy: serde_json::Value = serde_json::from_slice(&legacy_data).expect("legacy journal should decode as JSON");
assert_eq!(legacy["version"], serde_json::json!(1));
assert_eq!(legacy["backend_identity"], serde_json::Value::Null);
}
async fn put_real_cold_fill_object(store: &Arc<ECStore>, bucket: &str, object: &str, body: &[u8]) -> ObjectInfo {
let mut reader = PutObjReader::from_vec(body.to_vec());
store
+7
View File
@@ -390,6 +390,13 @@ pub(crate) mod bucket {
pub(crate) mod tier_delete_journal {
use std::sync::Arc;
pub(crate) fn record_tier_delete_journal_backend_identity(
je: &mut super::tier_sweeper::Jentry,
metadata: &std::collections::HashMap<String, String>,
) -> std::io::Result<()> {
crate::storage::storage_api::ecstore_bucket::lifecycle::tier_delete_journal::record_tier_delete_journal_backend_identity(je, metadata)
}
pub(crate) async fn persist_tier_delete_journal_entry(
api: Arc<crate::storage::storage_api::ECStore>,
je: &super::tier_sweeper::Jentry,
+4 -2
View File
@@ -492,7 +492,7 @@ pub(crate) mod ecstore_storage {
}
pub(crate) mod ecstore_tier {
pub(crate) use rustfs_ecstore::api::tier::tier::TierConfigMgr;
pub(crate) use rustfs_ecstore::api::tier::tier::{TierConfigMgr, TierConfigUpdateError};
pub(crate) use rustfs_ecstore::api::tier::{tier, tier_admin, tier_config, tier_handlers};
// Shared lifecycle/tier test utilities behind ecstore's `test-util` feature
// (rustfs/backlog#1148 ilm-6). Only linked into test builds.
@@ -577,6 +577,7 @@ pub(crate) type ReplicationStatusType = ecstore_bucket::replication::Replication
pub(crate) type ReplicationStats = StorageReplicationStatsHandle;
pub(crate) type StorageError = ecstore_error::StorageError;
pub(crate) type TierConfigMgr = ecstore_tier::TierConfigMgr;
pub(crate) type TierConfigUpdateError = ecstore_tier::TierConfigUpdateError;
pub(crate) use ecstore_disk::validate_batch_read_version_item_count;
pub(crate) type TransitionState = ecstore_bucket::lifecycle::bucket_lifecycle_ops::TransitionState;
pub(crate) type Error = ecstore_error::Error;
@@ -1438,7 +1439,8 @@ pub(crate) fn topology_snapshot_from_endpoint_pools_with_capabilities(
}
pub(crate) async fn reload_transition_tier_config(api: Arc<ECStore>) -> std::io::Result<()> {
ecstore_runtime::global_tier_config_mgr().write().await.reload(api).await
let handle = api.tier_config_mgr();
TierConfigMgr::reload_handle(&handle, api).await
}
pub(crate) async fn all_local_disk_path() -> Vec<String> {
@@ -16,8 +16,8 @@ readonly REQUESTS=2000
readonly KEY_MATRIX="1 4 32"
readonly REQUIRED_MEMORY_BYTES=$((28 * 1024 * 1024 * 1024))
readonly OBJECT_SIZE=$((384 * 1024 * 1024))
readonly READER_BYTES_QUERY='sum(rustfs_io_get_object_reader_bytes_total) or vector(0)'
readonly FOLLOWER_PERMIT_QUERY='sum(rustfs_object_data_cache_cold_fill_follower_disk_permits)'
readonly READER_BYTES_METRIC='rustfs_io_get_object_reader_bytes_total'
readonly FOLLOWER_PERMIT_METRIC='rustfs_object_data_cache_cold_fill_follower_disk_permits'
RUN=0
SELF_TEST=0
@@ -35,6 +35,7 @@ RESET_COMMAND=""
CACHE_OFF_COMMAND=""
CACHE_ON_COMMAND=""
SERVER_PID=""
SERVER_PID_COMMAND=""
CGROUP_PATH=""
SAMPLE_INTERVAL="0.25"
METRICS_SETTLE_SECONDS="5"
@@ -49,7 +50,8 @@ Usage:
--expected-sha256 HEX --prometheus-query-url URL \
--prometheus-scrape-seconds SECONDS \
--reset-command COMMAND --cache-off-command COMMAND \
--cache-on-command COMMAND --server-pid PID [--cgroup-path PATH] [options]
--cache-on-command COMMAND \
(--server-pid PID | --server-pid-command COMMAND) [--cgroup-path PATH] [options]
Required run contract:
* N is fixed at 2000; K is fixed at 1, 4, and 32; each K runs >= 3 rounds.
@@ -59,17 +61,17 @@ Required run contract:
* The RustFS process must be isolated from unrelated traffic and run in a
dedicated cgroup v2 with memory.max exactly 28 GiB.
* The built-in first-party reader-byte counter and follower-permit gauge must
each return exactly one Prometheus vector sample.
* Cache switch commands must return only after the mode is effective and must
keep the same RustFS PID inside the dedicated cgroup.
each return one value vector and one raw-scrape timestamp vector.
* Cache switch commands must return only after the mode is effective. A dynamic
PID resolver may observe a restarted RustFS process, but every process must
remain inside the same dedicated cgroup.
* Every response must match SHA-256, ETag, Content-Length, and its per-key
stable header contract. Date, x-amz-id-2, x-amz-request-id,
x-minio-request-id, x-rustfs-request-id, Connection, Keep-Alive, and
Transfer-Encoding are explicitly excluded as volatile/transport headers.
Built-in PromQL:
reader bytes: sum(rustfs_io_get_object_reader_bytes_total) or vector(0)
follower permits: sum(rustfs_object_data_cache_cold_fill_follower_disk_permits)
Built-in PromQL observes each metric value together with max(timestamp(metric));
the HTTP API evaluation timestamp is never treated as a scrape timestamp.
Options:
--self-test Validate the follower sample append/check chain locally
@@ -79,6 +81,7 @@ Options:
--prometheus-scrape-seconds N Configured scrape interval; required
--metrics-settle-seconds N Wait for final Prometheus scrape (default: 5)
--load-timeout-seconds N Per-round load timeout (default: 1800)
--server-pid-command COMMAND Resolve the live RustFS PID after every operator command
--out-dir PATH Artifact directory (default: temporary)
--strict Missing prerequisite is FAIL instead of SKIP
-h, --help Show this help
@@ -135,15 +138,76 @@ run_operator_command() {
bash -c "$command"
}
parse_server_pid() {
local output=$1
[[ $output =~ ^[[:space:]]*([0-9]+)[[:space:]]*$ ]] || return 1
printf '%s\n' "${BASH_REMATCH[1]}"
}
resolve_server_pid() {
local output
if [[ -n $SERVER_PID_COMMAND ]]; then
output=$(env -u AWS_ACCESS_KEY_ID -u AWS_SECRET_ACCESS_KEY -u AWS_SESSION_TOKEN bash -c "$SERVER_PID_COMMAND") || return 1
parse_server_pid "$output"
else
parse_server_pid "$SERVER_PID"
fi
}
parse_prometheus_observation() {
python3 -c '
import json, math, sys
doc = json.load(sys.stdin)
if doc.get("status") != "success":
raise SystemExit("Prometheus query was not successful")
result = doc.get("data", {}).get("result", [])
if len(result) != 2:
raise SystemExit(f"expected value and sample_timestamp vectors, got {len(result)}")
observations = {}
for item in result:
kind = item.get("metric", {}).get("__rustfs_probe_kind")
pair = item.get("value")
if kind not in ("value", "sample_timestamp") or not isinstance(pair, list) or len(pair) != 2:
raise SystemExit("invalid Prometheus observation shape")
evaluation_timestamp = float(pair[0])
observed_value = float(pair[1])
if not math.isfinite(evaluation_timestamp) or not math.isfinite(observed_value):
raise SystemExit("Prometheus observation is not finite")
if kind in observations:
raise SystemExit(f"duplicate Prometheus observation kind: {kind}")
observations[kind] = observed_value
if set(observations) != {"value", "sample_timestamp"}:
raise SystemExit("Prometheus observation is incomplete")
if observations["value"] < 0 or observations["sample_timestamp"] <= 0:
raise SystemExit("Prometheus observation is outside the accepted range")
print(format(observations["sample_timestamp"], ".17g"), format(observations["value"], ".17g"))
'
}
follower_samples_self_test() (
command -v python3 >/dev/null 2>&1 || fail "python3 is required for --self-test"
local temp_dir samples_file ready_file result sample
local temp_dir samples_file ready_file result sample stale_response fresh_response
temp_dir=$(mktemp -d "${TMPDIR:-/tmp}/rustfs-follower-samples.XXXXXX")
trap 'rm -rf -- "$temp_dir"' EXIT
samples_file="$temp_dir/follower.samples"
ready_file="$temp_dir/ready"
printf '999\n' >"$ready_file"
stale_response='{"status":"success","data":{"result":[{"metric":{"__rustfs_probe_kind":"value"},"value":[2000,"0"]},{"metric":{"__rustfs_probe_kind":"sample_timestamp"},"value":[2000,"900"]}]}}'
sample=$(parse_prometheus_observation <<<"$stale_response") || fail "self-test could not parse stale Prometheus observation"
[[ $sample == '900 0' ]] || fail "Prometheus evaluation timestamp was mistaken for scrape timestamp: $sample"
fresh_response='{"status":"success","data":{"result":[{"metric":{"__rustfs_probe_kind":"sample_timestamp"},"value":[2000,"1000"]},{"metric":{"__rustfs_probe_kind":"value"},"value":[2000,"0"]}]}}'
sample=$(parse_prometheus_observation <<<"$fresh_response") || fail "self-test could not parse fresh Prometheus observation"
[[ $sample == '1000 0' ]] || fail "Prometheus scrape timestamp was not preserved: $sample"
[[ $(parse_server_pid '101') == 101 ]] || fail "single-line PID output was not accepted"
if parse_server_pid '101 102' >/dev/null 2>&1; then
fail "multi-value PID output unexpectedly passed"
fi
if parse_server_pid $'101\n102' >/dev/null 2>&1; then
fail "multi-line PID output unexpectedly passed"
fi
: >"$samples_file"
for sample in '1000 0' '1000 0' '1000 0'; do
append_follower_sample "$samples_file" "$sample" || fail "self-test could not append duplicate follower sample"
@@ -197,6 +261,7 @@ while (($#)); do
--cache-off-command) need_value "$@"; CACHE_OFF_COMMAND=$2; shift ;;
--cache-on-command) need_value "$@"; CACHE_ON_COMMAND=$2; shift ;;
--server-pid) need_value "$@"; SERVER_PID=$2; shift ;;
--server-pid-command) need_value "$@"; SERVER_PID_COMMAND=$2; shift ;;
--cgroup-path) need_value "$@"; CGROUP_PATH=$2; shift ;;
--rounds) need_value "$@"; ROUNDS=$2; shift ;;
--sample-interval) need_value "$@"; SAMPLE_INTERVAL=$2; shift ;;
@@ -251,7 +316,13 @@ import sys
raise SystemExit(0 if float(sys.argv[1]) >= float(sys.argv[2]) else 1)
PY
[[ $SERVER_PID =~ ^[0-9]+$ ]] || skip_or_fail "--server-pid is required to bind OOM evidence to RustFS"
if [[ -n $SERVER_PID && -n $SERVER_PID_COMMAND ]]; then
fail "--server-pid and --server-pid-command are mutually exclusive"
fi
if [[ -z $SERVER_PID && -z $SERVER_PID_COMMAND ]]; then
skip_or_fail "--server-pid or --server-pid-command is required to bind OOM evidence to RustFS"
fi
SERVER_PID=$(resolve_server_pid) || skip_or_fail "could not resolve exactly one RustFS PID"
[[ -r /proc/$SERVER_PID/cgroup ]] || skip_or_fail "cannot read /proc/$SERVER_PID/cgroup"
if [[ -z $CGROUP_PATH ]]; then
cgroup_relative=$(awk -F: '$1 == "0" { print $3; exit }' "/proc/$SERVER_PID/cgroup")
@@ -272,6 +343,18 @@ server_in_cgroup() {
}
server_in_cgroup || skip_or_fail "RustFS PID $SERVER_PID is not in $CGROUP_PATH"
refresh_server_pid() {
local previous_pid=$SERVER_PID
local resolved
resolved=$(resolve_server_pid) || fail "could not resolve exactly one live RustFS PID"
[[ -r /proc/$resolved/cgroup ]] || fail "cannot read /proc/$resolved/cgroup"
SERVER_PID=$resolved
server_in_cgroup || fail "RustFS PID $SERVER_PID is not in $CGROUP_PATH"
if [[ $previous_pid != "$SERVER_PID" ]] && grep -Fxq -- "$previous_pid" "$CGROUP_PATH/cgroup.procs"; then
fail "replaced RustFS PID $previous_pid is still running in $CGROUP_PATH"
fi
}
nofile_limit=$(ulimit -n)
[[ $nofile_limit =~ ^[0-9]+$ ]] || skip_or_fail "unable to determine the open-file limit"
((nofile_limit >= REQUESTS + 256)) || skip_or_fail "open-file limit must be at least $((REQUESTS + 256)), got $nofile_limit"
@@ -294,31 +377,39 @@ cleanup_load() {
trap cleanup_load EXIT
prometheus_sample() {
local query=$1
local metric=$1
local query
local response
query="label_replace(sum($metric), \"__rustfs_probe_kind\", \"value\", \"__name__\", \".*\") or label_replace(max(timestamp($metric)), \"__rustfs_probe_kind\", \"sample_timestamp\", \"__name__\", \".*\")"
response=$(curl --fail --silent --show-error --get \
--data-urlencode "query=$query" "$PROMETHEUS_QUERY_URL") || return 1
python3 -c '
import json, math, sys
doc = json.load(sys.stdin)
if doc.get("status") != "success":
raise SystemExit("Prometheus query was not successful")
result = doc.get("data", {}).get("result", [])
if len(result) != 1 or "value" not in result[0]:
raise SystemExit(f"expected exactly one vector sample, got {len(result)}")
timestamp = float(result[0]["value"][0])
value = float(result[0]["value"][1])
if not math.isfinite(timestamp) or not math.isfinite(value) or value < 0:
raise SystemExit(f"invalid metric value: {value}")
print(format(timestamp, ".17g"), format(value, ".17g"))
' <<<"$response"
parse_prometheus_observation <<<"$response"
}
prometheus_value() {
prometheus_value_after() {
local metric=$1
local minimum_epoch=$2
local timestamp value
read -r timestamp value < <(prometheus_sample "$1") || return 1
[[ -n $timestamp && -n $value ]] || return 1
printf '%s\n' "$value"
local attempts
attempts=$(python3 - "$METRICS_SETTLE_SECONDS" "$PROMETHEUS_SCRAPE_SECONDS" "$SAMPLE_INTERVAL" <<'PY'
import math, sys
settle, scrape, poll = map(float, sys.argv[1:])
print(max(1, math.ceil((settle + 2 * scrape) / poll)))
PY
) || return 1
for ((attempt = 0; attempt < attempts; attempt++)); do
if read -r timestamp value < <(prometheus_sample "$metric") \
&& python3 - "$timestamp" "$minimum_epoch" <<'PY'
import sys
raise SystemExit(0 if float(sys.argv[1]) >= float(sys.argv[2]) else 1)
PY
then
printf '%s\n' "$value"
return 0
fi
sleep "$SAMPLE_INTERVAL"
done
return 1
}
cgroup_event() {
@@ -565,6 +656,9 @@ switch_cache_mode() {
local key_count=$4
run_operator_command "$mode" "$round" "$key_count" "$command" \
|| fail "cache-$mode command failed for round=$round K=$key_count"
if [[ -n $SERVER_PID_COMMAND ]]; then
refresh_server_pid
fi
server_in_cgroup || fail "cache-$mode command moved RustFS PID $SERVER_PID out of $CGROUP_PATH"
}
@@ -574,6 +668,9 @@ reset_cache() {
local key_count=$3
run_operator_command "$mode" "$round" "$key_count" "$RESET_COMMAND" \
|| fail "reset command failed for mode=$mode round=$round K=$key_count"
if [[ -n $SERVER_PID_COMMAND ]]; then
refresh_server_pid
fi
server_in_cgroup || fail "reset command moved RustFS PID $SERVER_PID out of $CGROUP_PATH"
}
@@ -594,14 +691,15 @@ for ((round = 1; round <= ROUNDS; round++)); do
printf 'INFO: K=%d round=%d: switching cache off for one request per key baseline\n' "$key_count" "$round"
switch_cache_mode off "$CACHE_OFF_COMMAND" "$round" "$key_count"
reset_cache off "$round" "$key_count"
baseline_reader_before=$(prometheus_value "$READER_BYTES_QUERY") || fail "reader query failed before cache-off K=$key_count round=$round"
baseline_start_epoch=$(python3 -c 'import time; print(format(time.time(), ".17g"))')
baseline_reader_before=$(prometheus_value_after "$READER_BYTES_METRIC" "$baseline_start_epoch") || fail "fresh reader query failed before cache-off K=$key_count round=$round"
baseline_oom_before=$(cgroup_event oom) || fail "cannot read oom before baseline"
baseline_oom_kill_before=$(cgroup_event oom_kill) || fail "cannot read oom_kill before baseline"
baseline_peak_before=$(cat "$CGROUP_PATH/memory.peak" 2>/dev/null || printf 'NA')
run_load "$key_count" "$key_count" "$baseline_summary" "$baseline_ready"
sleep "$METRICS_SETTLE_SECONDS"
baseline_reader_after=$(prometheus_value "$READER_BYTES_QUERY") || fail "reader query failed after cache-off K=$key_count round=$round"
baseline_finished_epoch=$(python3 -c 'import time; print(format(time.time(), ".17g"))')
baseline_reader_after=$(prometheus_value_after "$READER_BYTES_METRIC" "$baseline_finished_epoch") || fail "fresh reader query failed after cache-off K=$key_count round=$round"
baseline_reader_delta=$(numeric_positive_delta "$baseline_reader_before" "$baseline_reader_after") \
|| fail "cache-off reader baseline is invalid for K=$key_count round=$round"
baseline_oom_after=$(cgroup_event oom) || fail "cannot read oom after baseline"
@@ -629,7 +727,8 @@ for ((round = 1; round <= ROUNDS; round++)); do
switch_cache_mode on "$CACHE_ON_COMMAND" "$round" "$key_count"
reset_cache on "$round" "$key_count"
reader_before=$(prometheus_value "$READER_BYTES_QUERY") || fail "reader query failed before K=$key_count round=$round"
round_start_epoch=$(python3 -c 'import time; print(format(time.time(), ".17g"))')
reader_before=$(prometheus_value_after "$READER_BYTES_METRIC" "$round_start_epoch") || fail "fresh reader query failed before K=$key_count round=$round"
oom_before=$(cgroup_event oom) || fail "cannot read oom before round"
oom_kill_before=$(cgroup_event oom_kill) || fail "cannot read oom_kill before round"
memory_peak_before=$(cat "$CGROUP_PATH/memory.peak" 2>/dev/null || printf 'NA')
@@ -639,7 +738,7 @@ for ((round = 1; round <= ROUNDS; round++)); do
LOAD_PID=$load_pid
while kill -0 "$load_pid" 2>/dev/null; do
if [[ -e $ready_file ]]; then
if sample=$(prometheus_sample "$FOLLOWER_PERMIT_QUERY" 2>/dev/null); then
if sample=$(prometheus_sample "$FOLLOWER_PERMIT_METRIC" 2>/dev/null); then
append_follower_sample "$follower_file" "$sample" || fail "invalid follower sample for K=$key_count round=$round"
fi
fi
@@ -648,8 +747,8 @@ for ((round = 1; round <= ROUNDS; round++)); do
wait "$load_pid" || fail "load or SHA-256 validation failed for K=$key_count round=$round (see $summary_file)"
LOAD_PID=""
sleep "$METRICS_SETTLE_SECONDS"
reader_after=$(prometheus_value "$READER_BYTES_QUERY") || fail "reader query failed after K=$key_count round=$round"
round_finished_epoch=$(python3 -c 'import time; print(format(time.time(), ".17g"))')
reader_after=$(prometheus_value_after "$READER_BYTES_METRIC" "$round_finished_epoch") || fail "fresh reader query failed after K=$key_count round=$round"
expected_reader_bytes=$(((baseline_reader_delta * 90) / 100))
reader_limit=$(((baseline_reader_delta * 110 + 99) / 100))
reader_delta=$(numeric_delta_check "$reader_before" "$reader_after" \