From 9f25858b05cbb16112fee243e7f71f2f8600a599 Mon Sep 17 00:00:00 2001 From: houseme Date: Tue, 21 Jul 2026 15:26:23 +0800 Subject: [PATCH] fix(tier): add mutation intent record store (#5075) Add canonical record-object naming and crate-private save, load, and idempotent delete helpers for tier mutation intents. Cover malformed record keys, mismatched persisted mutation IDs, and an ECStore config-object round trip under test-util. Co-authored-by: heihutu --- .../src/services/tier/tier_mutation_intent.rs | 197 ++++++++++++++++++ crates/ecstore/src/store/init.rs | 111 ++++++++++ 2 files changed, 308 insertions(+) diff --git a/crates/ecstore/src/services/tier/tier_mutation_intent.rs b/crates/ecstore/src/services/tier/tier_mutation_intent.rs index 82ad5f754..061617736 100644 --- a/crates/ecstore/src/services/tier/tier_mutation_intent.rs +++ b/crates/ecstore/src/services/tier/tier_mutation_intent.rs @@ -12,14 +12,22 @@ // See the License for the specific language governing permissions and // limitations under the License. +use std::sync::Arc; + use rustfs_utils::crypto::{hex_sha256, is_sha256_checksum}; use serde::{Deserialize, Serialize}; use uuid::Uuid; +use crate::config::com; +use crate::disk::RUSTFS_META_BUCKET; +use crate::error::{Error, Result as EcstoreResult}; use crate::services::tier::tier::TierDestinationId; +use crate::storage_api_contracts::list::ListOperations as _; +use crate::store::ECStore; pub(crate) const TIER_MUTATION_INTENT_SCHEMA: &str = "rustfs-tier-mutation-intent-v1"; pub(crate) const MAX_TIER_MUTATION_INTENT_SIZE: usize = 64 * 1024; +pub(crate) const TIER_MUTATION_INTENT_RECORD_PREFIX: &str = "tier/mutation-intents/records"; pub(crate) type TierMutationDigest = [u8; 32]; pub(crate) type Result = std::result::Result; @@ -102,6 +110,15 @@ pub(crate) struct TierMutationIntent { pub expires_at_unix_nanos: i64, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct TierMutationIntentRecordScan { + pub intents: Vec, + pub scanned: usize, + pub failed: usize, + pub next_marker: Option, + pub truncated: bool, +} + impl TierMutationIntent { pub(crate) fn validate(&self) -> Result<()> { if self.mutation_id.is_nil() { @@ -238,6 +255,127 @@ impl TierMutationIntent { } } +pub(crate) fn tier_mutation_intent_record_object_name(mutation_id: Uuid) -> Result { + if mutation_id.is_nil() { + return Err(TierMutationIntentError::Corrupt("mutation_id is nil")); + } + let mutation_key = mutation_id.simple().to_string(); + Ok(format!( + "{}/{}/{}/{}.json", + TIER_MUTATION_INTENT_RECORD_PREFIX, + &mutation_key[..2], + &mutation_key[2..4], + mutation_key + )) +} + +pub(crate) fn tier_mutation_intent_id_from_record_object_name(object: &str) -> Result { + let prefix = format!("{TIER_MUTATION_INTENT_RECORD_PREFIX}/"); + let suffix = object + .strip_prefix(&prefix) + .ok_or(TierMutationIntentError::Corrupt("intent record path has wrong prefix"))?; + let mut parts = suffix.split('/'); + let shard_a = parts + .next() + .ok_or(TierMutationIntentError::Corrupt("intent record path is incomplete"))?; + let shard_b = parts + .next() + .ok_or(TierMutationIntentError::Corrupt("intent record path is incomplete"))?; + let file_name = parts + .next() + .ok_or(TierMutationIntentError::Corrupt("intent record path is incomplete"))?; + if parts.next().is_some() { + return Err(TierMutationIntentError::Corrupt("intent record path is not canonical")); + } + let mutation_key = file_name + .strip_suffix(".json") + .ok_or(TierMutationIntentError::Corrupt("intent record path has wrong suffix"))?; + if mutation_key.len() != 32 + || !mutation_key + .bytes() + .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase()) + { + return Err(TierMutationIntentError::Corrupt("intent record path has invalid mutation id")); + } + if shard_a != &mutation_key[..2] || shard_b != &mutation_key[2..4] { + return Err(TierMutationIntentError::Corrupt("intent record path shard does not match mutation id")); + } + Uuid::parse_str(mutation_key).map_err(|_| TierMutationIntentError::Corrupt("intent record path has invalid uuid")) +} + +pub(crate) async fn save_tier_mutation_intent_record(api: Arc, intent: &TierMutationIntent) -> EcstoreResult<()> { + let object = tier_mutation_intent_record_object_name(intent.mutation_id).map_err(tier_mutation_intent_store_error)?; + let data = intent.encode().map_err(tier_mutation_intent_store_error)?; + com::save_config(api, &object, data).await +} + +pub(crate) async fn load_tier_mutation_intent_record(api: Arc, mutation_id: Uuid) -> EcstoreResult { + let object = tier_mutation_intent_record_object_name(mutation_id).map_err(tier_mutation_intent_store_error)?; + let data = com::read_config(api, &object).await?; + TierMutationIntent::decode(mutation_id, &data).map_err(tier_mutation_intent_store_error) +} + +pub(crate) async fn delete_tier_mutation_intent_record(api: Arc, mutation_id: Uuid) -> EcstoreResult<()> { + let object = tier_mutation_intent_record_object_name(mutation_id).map_err(tier_mutation_intent_store_error)?; + match com::delete_config(api, &object).await { + Ok(()) | Err(Error::ConfigNotFound) => Ok(()), + Err(err) => Err(err), + } +} + +pub(crate) async fn list_tier_mutation_intent_records( + api: Arc, + limit: usize, + marker: Option, +) -> EcstoreResult { + if limit == 0 { + return Err(Error::other("tier mutation intent scan limit must be greater than zero")); + } + + let list = api + .clone() + .list_objects_v2( + RUSTFS_META_BUCKET, + TIER_MUTATION_INTENT_RECORD_PREFIX, + marker, + None, + i32::try_from(limit).map_or(i32::MAX, |value| value), + false, + None, + false, + ) + .await?; + let mut scan = TierMutationIntentRecordScan { + intents: Vec::with_capacity(list.objects.len()), + scanned: 0, + failed: 0, + next_marker: list.next_continuation_token, + truncated: list.is_truncated, + }; + + for object in list.objects { + scan.scanned += 1; + let mutation_id = match tier_mutation_intent_id_from_record_object_name(&object.name) { + Ok(mutation_id) => mutation_id, + Err(_) => { + scan.failed += 1; + continue; + } + }; + match load_tier_mutation_intent_record(api.clone(), mutation_id).await { + Ok(intent) => scan.intents.push(intent), + Err(Error::ConfigNotFound) => {} + Err(_) => scan.failed += 1, + } + } + + Ok(scan) +} + +fn tier_mutation_intent_store_error(err: TierMutationIntentError) -> Error { + Error::other(err) +} + #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(deny_unknown_fields)] struct PersistedTierMutationIntent { @@ -384,6 +522,65 @@ mod tests { assert_eq!(decoded.kind, TierMutationIntentKind::Add); } + #[test] + fn intent_record_object_name_is_canonical_and_reversible() { + let mutation_id = Uuid::parse_str("12345678-1234-5678-9abc-def012345678").expect("test uuid should parse"); + + let object = tier_mutation_intent_record_object_name(mutation_id).expect("record object should build"); + let parsed = tier_mutation_intent_id_from_record_object_name(&object).expect("record object should parse"); + + assert_eq!(object, "tier/mutation-intents/records/12/34/12345678123456789abcdef012345678.json"); + assert_eq!(parsed, mutation_id); + } + + #[test] + fn intent_record_object_name_rejects_nil_and_malformed_paths() { + assert!(matches!( + tier_mutation_intent_record_object_name(Uuid::nil()), + Err(TierMutationIntentError::Corrupt("mutation_id is nil")) + )); + assert!(matches!( + tier_mutation_intent_id_from_record_object_name("tier/mutation-intents/records/12/34/not-a-uuid.json"), + Err(TierMutationIntentError::Corrupt("intent record path has invalid mutation id")) + )); + assert!(matches!( + tier_mutation_intent_id_from_record_object_name( + "tier/mutation-intents/records/12/34/12345678123456789abcdef012345678" + ), + Err(TierMutationIntentError::Corrupt("intent record path has wrong suffix")) + )); + assert!(matches!( + tier_mutation_intent_id_from_record_object_name( + "tier/mutation-intents/records/12/35/12345678123456789abcdef012345678.json" + ), + Err(TierMutationIntentError::Corrupt("intent record path shard does not match mutation id")) + )); + assert!(matches!( + tier_mutation_intent_id_from_record_object_name( + "tier/mutation-intents/records/12/34/12345678123456789abcdef012345678/extra.json" + ), + Err(TierMutationIntentError::Corrupt("intent record path is not canonical")) + )); + assert!(matches!( + tier_mutation_intent_id_from_record_object_name( + "tier/mutation-intents/records/12/34/12345678123456789ABCDEF012345678.json" + ), + Err(TierMutationIntentError::Corrupt("intent record path has invalid mutation id")) + )); + } + + #[test] + fn intent_decode_rejects_record_key_mismatch() { + let intent = prepared_intent(); + let encoded = intent.encode().expect("intent should encode"); + let wrong_id = Uuid::new_v4(); + + assert!(matches!( + TierMutationIntent::decode(wrong_id, &encoded), + Err(TierMutationIntentError::Corrupt("mutation_id does not match intent key")) + )); + } + #[test] fn intent_validation_rejects_target_shape_mismatch() { let mut add = prepared_intent(); diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 9a11e3174..81af5fc3f 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -541,11 +541,17 @@ mod tests { }, }, client::transition_api::ReaderImpl, + config::com, disk::RUSTFS_META_BUCKET, runtime::{global::set_object_store_resolver, sources as runtime_sources}, services::tier::{ test_util::{MockWarmBackend, TransitionCleanupStoreBarrier, register_mock_tier}, tier::TierConfigMgr, + tier_mutation_intent::{ + TIER_MUTATION_INTENT_RECORD_PREFIX, TierMutationIntent, TierMutationIntentKind, TierMutationIntentState, + TierMutationIntentTarget, delete_tier_mutation_intent_record, list_tier_mutation_intent_records, + load_tier_mutation_intent_record, save_tier_mutation_intent_record, + }, warm_backend::WarmBackend, }, storage_api_contracts::{ @@ -1229,6 +1235,111 @@ mod tests { shutdown_b.cancel(); } + #[cfg(feature = "test-util")] + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn tier_mutation_intent_record_round_trips_through_config_store() { + let temp_dir = tempfile::tempdir().expect("create temp store dir"); + let (_ctx, store, _shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "tier-mutation-intent-record", &[4])).await; + let mutation_id = uuid::Uuid::new_v4(); + let intent = TierMutationIntent { + mutation_id, + revision: 1, + kind: TierMutationIntentKind::Edit, + state: TierMutationIntentState::Prepared, + old_config_etag: Some("old-etag".to_string()), + committed_config_etag: None, + candidate_digest: [3; 32], + affected_targets: vec![TierMutationIntentTarget { + tier_name: "COLD-A".to_string(), + old_backend_identity: Some([1; 32]), + new_backend_identity: Some([2; 32]), + }], + expires_at_unix_nanos: 1_780_000_000_000_000_000, + }; + + save_tier_mutation_intent_record(store.clone(), &intent) + .await + .expect("tier mutation intent record should persist"); + let loaded = load_tier_mutation_intent_record(store.clone(), mutation_id) + .await + .expect("tier mutation intent record should load"); + + assert_eq!(loaded, intent); + + delete_tier_mutation_intent_record(store.clone(), mutation_id) + .await + .expect("tier mutation intent record delete should be idempotent"); + delete_tier_mutation_intent_record(store.clone(), mutation_id) + .await + .expect("tier mutation intent record delete should tolerate missing records"); + let err = load_tier_mutation_intent_record(store, mutation_id) + .await + .expect_err("deleted tier mutation intent record should not load"); + assert!(matches!(err, Error::ConfigNotFound)); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn tier_mutation_intent_record_scan_retains_good_records_and_counts_bad_records() { + let temp_dir = tempfile::tempdir().expect("create temp store dir"); + let (_ctx, store, _shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "tier-mutation-intent-scan", &[4])).await; + let build_intent = |mutation_id: uuid::Uuid, tier_name: &str| TierMutationIntent { + mutation_id, + revision: 1, + kind: TierMutationIntentKind::Edit, + state: TierMutationIntentState::Prepared, + old_config_etag: Some("old-etag".to_string()), + committed_config_etag: None, + candidate_digest: [3; 32], + affected_targets: vec![TierMutationIntentTarget { + tier_name: tier_name.to_string(), + old_backend_identity: Some([1; 32]), + new_backend_identity: Some([2; 32]), + }], + expires_at_unix_nanos: 1_780_000_000_000_000_000, + }; + let first_id = uuid::Uuid::parse_str("12345678-1234-5678-9abc-def012345678").expect("first uuid should parse"); + let second_id = uuid::Uuid::parse_str("22345678-1234-5678-9abc-def012345678").expect("second uuid should parse"); + let first = build_intent(first_id, "COLD-A"); + let second = build_intent(second_id, "COLD-B"); + save_tier_mutation_intent_record(store.clone(), &first) + .await + .expect("first tier mutation intent record should persist"); + save_tier_mutation_intent_record(store.clone(), &second) + .await + .expect("second tier mutation intent record should persist"); + com::save_config( + store.clone(), + &format!("{TIER_MUTATION_INTENT_RECORD_PREFIX}/00/00/33345678123456789abcdef012345678.json"), + b"{}".to_vec(), + ) + .await + .expect("malformed-shard intent record should persist"); + com::save_config( + store.clone(), + &format!("{TIER_MUTATION_INTENT_RECORD_PREFIX}/44/44/44445678123456789abcdef012345678.json"), + b"{".to_vec(), + ) + .await + .expect("corrupt-json intent record should persist"); + + let scan = list_tier_mutation_intent_records(store, 100, None) + .await + .expect("tier mutation intent records should scan"); + let mut loaded_ids: Vec<_> = scan.intents.into_iter().map(|intent| intent.mutation_id).collect(); + loaded_ids.sort(); + + assert_eq!(scan.scanned, 4); + assert_eq!(scan.failed, 2); + assert_eq!(loaded_ids, vec![first_id, second_id]); + assert!(!scan.truncated); + assert_eq!(scan.next_marker, None); + } + #[cfg(feature = "test-util")] #[tokio::test] #[serial_test::serial(storage_class_env)]