From 937b31131690c823ebbb7c41464dffbc95baa275 Mon Sep 17 00:00:00 2001 From: houseme Date: Tue, 21 Jul 2026 23:16:12 +0800 Subject: [PATCH] fix(tier): lock tier config mutations (#5080) * fix(tier): lock tier config mutations Co-Authored-By: heihutu * fix(tier): add mutation RPC auth contract (#5082) Co-authored-by: heihutu * fix(tier): add peer mutation handler core (#5084) Co-authored-by: heihutu * fix(tier): add mutation control rpc service (#5087) Co-authored-by: heihutu * fix(tier): recover prepared mutation drains (#5093) Recover prepared tier mutation intent records into the local tier runtime so a restarted peer fails closed before issuing new remote-tier operation leases or conflicting admin publishes. Reconcile the recovered block map on each scan so committed, aborted, or removed intents clear stale local blocks instead of wedging the peer until process restart. Co-authored-by: heihutu * fix(tier): prove zero references before tier removal (#5092) Signed-off-by: houseme Co-authored-by: heihutu * fix(tier): clear peer mutation runtime blocks (#5094) Install prepared mutation runtime blocks when peer prepare requests are applied or replayed so followers fail closed immediately before restart recovery. Clear the in-memory block once peer commit or abort reaches a durable terminal state, including delayed duplicate prepare requests that observe a committed or aborted record. Co-authored-by: heihutu --------- Signed-off-by: houseme Co-authored-by: heihutu --- crates/ecstore/src/api/mod.rs | 7 + crates/ecstore/src/cluster/rpc/http_auth.rs | 50 + crates/ecstore/src/services/tier/mod.rs | 1 + crates/ecstore/src/services/tier/tier.rs | 1275 ++++++++++++++++- .../src/services/tier/tier_mutation_intent.rs | 24 +- .../src/services/tier/tier_mutation_peer.rs | 278 ++++ crates/ecstore/src/store/init.rs | 222 +++ .../src/generated/proto_gen/node_service.rs | 403 ++++++ crates/protos/src/lib.rs | 271 ++++ crates/protos/src/node.proto | 39 + rustfs/src/server/http.rs | 53 +- rustfs/src/storage/rpc/mod.rs | 5 +- rustfs/src/storage/rpc/node_service.rs | 358 ++++- rustfs/src/storage/storage_api.rs | 6 +- rustfs/src/storage/tonic_service.rs | 2 +- rustfs/src/storage_api.rs | 2 +- 16 files changed, 2947 insertions(+), 49 deletions(-) create mode 100644 crates/ecstore/src/services/tier/tier_mutation_peer.rs diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index 31a34e37d..da875822e 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -437,6 +437,13 @@ pub mod tier { }; } + pub mod tier_mutation_peer { + pub use crate::services::tier::tier_mutation_peer::{ + MAX_TIER_MUTATION_PEER_COMMIT_ETAG_SIZE, TierMutationPeerError, TierMutationPeerOutcome, TierMutationPeerResult, + TierMutationPeerState, handle_tier_mutation_peer_request, + }; + } + pub mod warm_backend { pub use crate::services::tier::warm_backend::{ WarmBackend, WarmBackendGetOpts, WarmBackendImpl, build_transition_put_options, check_warm_backend, new_warm_backend, diff --git a/crates/ecstore/src/cluster/rpc/http_auth.rs b/crates/ecstore/src/cluster/rpc/http_auth.rs index 4f479101f..1b1802795 100644 --- a/crates/ecstore/src/cluster/rpc/http_auth.rs +++ b/crates/ecstore/src/cluster/rpc/http_auth.rs @@ -1144,6 +1144,56 @@ mod tests { assert_eq!(tampered.to_string(), "RPC content SHA-256 mismatch"); } + #[test] + fn tier_mutation_rpc_contract_requires_method_bound_v2_body_digest() { + ensure_test_rpc_secret(); + let mutation_id = uuid::uuid!("12345678-1234-5678-9abc-def012345678"); + let body = rustfs_protos::canonical_tier_mutation_rpc_body( + rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION, + rustfs_protos::TierMutationRpcPhase::Prepare, + mutation_id, + b"canonical-tier-mutation-prepare", + ) + .expect("small tier mutation body should encode"); + let mut request = tonic::Request::new(()); + set_tonic_canonical_body_digest(&mut request, &body).expect("canonical body digest should be attached"); + let content_sha256 = request + .metadata() + .get(RPC_CONTENT_SHA256_HEADER) + .and_then(|value| value.to_str().ok()); + let headers = gen_tonic_signature_headers( + "node-a:9000", + "node_service.TierMutationControlService", + "PrepareTierMutation", + content_sha256, + ) + .expect("body-bound tier mutation auth headers should build"); + request.metadata_mut().as_mut().extend(headers.clone()); + + assert!( + verify_tonic_rpc_signature("node-a:9000", "/node_service.TierMutationControlService/PrepareTierMutation", &headers) + .is_ok(), + "tier mutation RPC signature must bind destination, service, method, nonce, and body digest" + ); + let method_replay = + verify_tonic_rpc_signature("node-a:9000", "/node_service.TierMutationControlService/CommitTierMutation", &headers) + .expect_err("prepare auth must not replay to commit"); + assert_eq!(method_replay.to_string(), "Invalid RPC v2 signature"); + let service_replay = verify_tonic_rpc_signature("node-a:9000", "/node_service.NodeService/PrepareTierMutation", &headers) + .expect_err("tier mutation auth must not replay to the legacy node service path"); + assert_eq!(service_replay.to_string(), "Invalid RPC v2 signature"); + let tampered_body = rustfs_protos::canonical_tier_mutation_rpc_body( + rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION, + rustfs_protos::TierMutationRpcPhase::Commit, + mutation_id, + b"canonical-tier-mutation-prepare", + ) + .expect("small tier mutation body should encode"); + let tampered = + verify_tonic_canonical_body_digest(&request, &tampered_body).expect_err("commit body must not match prepare digest"); + assert_eq!(tampered.to_string(), "RPC content SHA-256 mismatch"); + } + #[test] fn partial_v2_metadata_fails_closed() { ensure_test_rpc_secret(); diff --git a/crates/ecstore/src/services/tier/mod.rs b/crates/ecstore/src/services/tier/mod.rs index 97b6ea6be..acc0bd93b 100644 --- a/crates/ecstore/src/services/tier/mod.rs +++ b/crates/ecstore/src/services/tier/mod.rs @@ -20,6 +20,7 @@ pub mod tier_config; pub mod tier_gen; pub mod tier_handlers; pub(crate) mod tier_mutation_intent; +pub mod tier_mutation_peer; pub mod warm_backend; pub mod warm_backend_aliyun; pub mod warm_backend_azure; diff --git a/crates/ecstore/src/services/tier/tier.rs b/crates/ecstore/src/services/tier/tier.rs index 26bb97e85..38a4267a9 100644 --- a/crates/ecstore/src/services/tier/tier.rs +++ b/crates/ecstore/src/services/tier/tier.rs @@ -56,6 +56,9 @@ use crate::services::tier::{ warm_backend::{WarmBackend, check_warm_backend, new_warm_backend}, }; use crate::storage_api_contracts::{ + bucket::BucketOperations, + list::{ListOperations, StorageListObjectVersionsInfo, StorageListObjectsV2Info, StorageObjectInfoOrErr, StorageWalkOptions}, + namespace::NamespaceLocking, object::{ DeletedObject, EcstoreObjectIO, EcstoreObjectOperations, HTTPPreconditions, ObjectIO, ObjectOperations, ObjectToDelete, }, @@ -66,6 +69,7 @@ use crate::{ disk::{MIGRATING_META_BUCKET, RUSTFS_META_BUCKET}, object_api::{GetObjectReader, ObjectInfo, ObjectOptions, PutObjReader}, runtime::sources as runtime_sources, + set_disk::get_lock_acquire_timeout, store::ECStore, }; use rustfs_filemeta::FileInfo; @@ -75,12 +79,21 @@ use s3s::S3ErrorCode; use super::{ tier_handlers::{ERR_TIER_BUCKET_NOT_FOUND, ERR_TIER_CONNECT_ERR, ERR_TIER_INVALID_CREDENTIALS, ERR_TIER_PERM_ERR}, + tier_mutation_intent::{ + TierMutationIntent, TierMutationIntentKind, TierMutationIntentState, TierMutationIntentTarget, + list_tier_mutation_intent_records, + }, warm_backend::WarmBackendImpl, }; const TIER_CFG_REFRESH: Duration = Duration::from_secs(15 * 60); const TIER_OPERATION_DRAIN_TIMEOUT: Duration = Duration::from_secs(30); const TIER_REMOTE_VALIDATION_TIMEOUT: Duration = Duration::from_secs(30); +const TIER_REFERENCE_PROOF_LIST_LIMIT: i32 = 1000; +const TIER_MUTATION_INTENT_RECOVERY_SCAN_LIMIT: usize = 1000; +const TIER_MUTATION_BLOCK_MESSAGE: &str = "Remote tier configuration is being replaced"; + +type TierReferenceProofWalkOptions = StorageWalkOptions bool>; fn delayed_tier_refresh_interval(period: Duration) -> tokio::time::Interval { interval_at(Instant::now() + period, period) @@ -193,6 +206,7 @@ type SharedWarmBackend = Arc; struct TierDriverRuntime { generations: HashMap>, draining: HashMap, + prepared_mutation_blocks: HashMap, next_generation: u64, next_drain_epoch: u64, admin_updates: Arc>, @@ -366,6 +380,14 @@ fn registered_tier_driver_runtime(manager: &TierConfigMgr) -> Option bool { + runtime.draining.contains_key(tier_name) || runtime.prepared_mutation_blocks.contains_key(tier_name) +} + +fn runtime_has_mutation_block(runtime: &TierDriverRuntime) -> bool { + !runtime.draining.is_empty() || !runtime.prepared_mutation_blocks.is_empty() +} + type TierDriverFingerprint = [u8; 32]; pub(crate) type TierDestinationId = [u8; 32]; pub(crate) type DriverRevision = u64; @@ -403,15 +425,37 @@ enum TierCandidateMutation { } impl TierCandidateMutation { + fn intent_kind(&self) -> TierMutationIntentKind { + match self { + Self::Add(_, _) => TierMutationIntentKind::Add, + Self::Edit(_, _) => TierMutationIntentKind::Edit, + Self::Remove(_, _) => TierMutationIntentKind::Remove, + Self::Clear(_) => TierMutationIntentKind::Clear, + } + } + + fn explicit_tier_name(&self) -> Option<&str> { + match self { + Self::Add(config, _) => Some(&config.name), + Self::Edit(tier_name, _) | Self::Remove(tier_name, _) => Some(tier_name), + Self::Clear(_) => None, + } + } + fn target_tiers(&self, manager: &TierConfigMgr, candidate: &TierConfigMgr) -> HashSet { let mut targets = changed_tier_names(manager, candidate); match self { Self::Add(config, _) => { targets.insert(config.name.clone()); } - Self::Edit(tier_name, _) | Self::Remove(tier_name, _) => { + Self::Edit(tier_name, _) => { targets.insert(tier_name.clone()); } + Self::Remove(tier_name, _) => { + if manager.tiers.contains_key(tier_name) || candidate.tiers.contains_key(tier_name) { + targets.insert(tier_name.clone()); + } + } Self::Clear(_) => { targets.extend(manager.tiers.keys().chain(candidate.tiers.keys()).cloned()); } @@ -419,6 +463,20 @@ impl TierCandidateMutation { targets } + fn affected_targets( + &self, + manager: &TierConfigMgr, + candidate: &TierConfigMgr, + ) -> std::result::Result, AdminError> { + let explicit_tier_name = self.explicit_tier_name(); + build_tier_mutation_affected_targets( + self.intent_kind(), + tier_mutation_proof_targets(self.intent_kind(), explicit_tier_name, manager, candidate), + manager, + candidate, + ) + } + async fn apply(self, candidate: &mut TierConfigMgr) -> std::result::Result, AdminError> { match self { Self::Add(config, force) => { @@ -442,6 +500,210 @@ impl TierCandidateMutation { } } +fn tier_mutation_proof_targets( + kind: TierMutationIntentKind, + explicit_tier_name: Option<&str>, + current: &TierConfigMgr, + candidate: &TierConfigMgr, +) -> HashSet { + let mut targets = changed_tier_names(current, candidate); + match kind { + TierMutationIntentKind::Add | TierMutationIntentKind::Edit | TierMutationIntentKind::Remove => { + if let Some(tier_name) = explicit_tier_name + && (current.tiers.contains_key(tier_name) || candidate.tiers.contains_key(tier_name)) + { + targets.insert(tier_name.to_string()); + } + } + TierMutationIntentKind::Clear => { + targets.extend(current.tiers.keys().chain(candidate.tiers.keys()).cloned()); + } + } + targets +} + +fn build_tier_mutation_affected_targets( + kind: TierMutationIntentKind, + target_tiers: HashSet, + current: &TierConfigMgr, + candidate: &TierConfigMgr, +) -> std::result::Result, AdminError> { + let mut targets = target_tiers.into_iter().collect::>(); + targets.sort(); + let mut out = Vec::with_capacity(targets.len()); + for tier_name in targets { + let old_backend_identity = current + .tiers + .get(&tier_name) + .map(tier_backend_identity) + .transpose() + .map_err(tier_backend_identity_admin_error)?; + let new_backend_identity = candidate + .tiers + .get(&tier_name) + .map(tier_backend_identity) + .transpose() + .map_err(tier_backend_identity_admin_error)?; + if old_backend_identity.is_none() && new_backend_identity.is_none() { + continue; + } + let target = TierMutationIntentTarget { + tier_name, + old_backend_identity, + new_backend_identity, + }; + validate_tier_mutation_target_shape(kind, &target)?; + out.push(target); + } + Ok(out) +} + +fn validate_tier_mutation_target_shape( + kind: TierMutationIntentKind, + target: &TierMutationIntentTarget, +) -> std::result::Result<(), AdminError> { + let valid = match kind { + TierMutationIntentKind::Add => target.old_backend_identity.is_none() && target.new_backend_identity.is_some(), + TierMutationIntentKind::Edit => target.old_backend_identity.is_some() && target.new_backend_identity.is_some(), + TierMutationIntentKind::Remove | TierMutationIntentKind::Clear => { + target.old_backend_identity.is_some() && target.new_backend_identity.is_none() + } + }; + if valid { + return Ok(()); + } + let mut err = ERR_TIER_INVALID_CONFIG.clone(); + err.message = "Remote tier mutation target identity proof does not match mutation kind".to_string(); + Err(err) +} + +fn tier_backend_identity_admin_error(err: io::Error) -> AdminError { + let mut admin_err = ERR_TIER_INVALID_CONFIG.clone(); + admin_err.message = err.to_string(); + admin_err +} + +trait TierReferenceProofStore: + EcstoreObjectIO + + BucketOperations + + ListOperations< + Error = Error, + ListObjectsV2Info = StorageListObjectsV2Info, + ListObjectVersionsInfo = StorageListObjectVersionsInfo, + ObjectInfoOrErr = StorageObjectInfoOrErr, + WalkOptions = TierReferenceProofWalkOptions, + WalkCancellation = tokio_util::sync::CancellationToken, + WalkResultSender = tokio::sync::mpsc::Sender>, + > +{ +} + +impl TierReferenceProofStore for T where + T: EcstoreObjectIO + + BucketOperations + + ListOperations< + Error = Error, + ListObjectsV2Info = StorageListObjectsV2Info, + ListObjectVersionsInfo = StorageListObjectVersionsInfo, + ObjectInfoOrErr = StorageObjectInfoOrErr, + WalkOptions = TierReferenceProofWalkOptions, + WalkCancellation = tokio_util::sync::CancellationToken, + WalkResultSender = tokio::sync::mpsc::Sender>, + > +{ +} + +async fn ensure_no_authoritative_tier_object_references( + api: Arc, + affected_targets: &[TierMutationIntentTarget], +) -> std::result::Result<(), AdminError> +where + S: TierReferenceProofStore, +{ + for target in affected_targets { + if target.old_backend_identity.is_some() { + ensure_no_authoritative_target_references(api.clone(), target).await?; + } + } + Ok(()) +} + +async fn ensure_no_authoritative_target_references( + api: Arc, + target: &TierMutationIntentTarget, +) -> std::result::Result<(), AdminError> +where + S: TierReferenceProofStore, +{ + let buckets = api + .list_bucket(&Default::default()) + .await + .map_err(tier_reference_proof_admin_error)?; + for bucket in buckets { + let mut marker = None; + let mut version_marker = None; + loop { + let page = api + .clone() + .list_object_versions( + &bucket.name, + "", + marker.take(), + version_marker.take(), + None, + TIER_REFERENCE_PROOF_LIST_LIMIT, + ) + .await + .map_err(tier_reference_proof_admin_error)?; + for object in &page.objects { + if tier_object_blocks_target_rebind(object, target).map_err(tier_reference_proof_admin_error)? { + return Err(tier_reference_proof_in_use_error(&target.tier_name, object)); + } + } + if !page.is_truncated { + break; + } + marker = page.next_marker; + version_marker = page.next_version_idmarker; + if marker.is_none() { + return Err(tier_reference_proof_admin_error(io::Error::other( + "tier reference proof listing is truncated without a next marker", + ))); + } + } + } + Ok(()) +} + +fn tier_object_blocks_target_rebind(object: &ObjectInfo, target: &TierMutationIntentTarget) -> io::Result { + if object.transitioned_object.status != rustfs_filemeta::TRANSITION_COMPLETE + || object.transitioned_object.tier != target.tier_name + { + return Ok(false); + } + let current_identity = tier_destination_id_from_metadata(&object.user_defined)?; + Ok(match (¤t_identity, &target.new_backend_identity) { + (_, None) => true, + (Some(current), Some(new)) => current != new, + (None, Some(_)) => true, + }) +} + +fn tier_reference_proof_in_use_error(tier_name: &str, object: &ObjectInfo) -> AdminError { + let mut err = ERR_TIER_BACKEND_IN_USE.clone(); + err.message = format!( + "Remote tier {tier_name} still has object references, for example {}/{}", + object.bucket, object.name + ); + err +} + +fn tier_reference_proof_admin_error(err: impl std::fmt::Display) -> AdminError { + let mut admin_err = ERR_TIER_INVALID_CONFIG.clone(); + admin_err.message = format!("Remote tier reference proof failed: {err}"); + admin_err +} + async fn apply_tier_candidate_mutation( mutation: TierCandidateMutation, candidate: &mut TierConfigMgr, @@ -983,6 +1245,10 @@ fn tier_config_path(file: &str) -> String { format!("{}{}{}", CONFIG_PREFIX, SLASH_SEPARATOR, file) } +fn tier_config_lock_path() -> String { + format!("{}.lock", tier_config_path(TIER_CONFIG_FILE)) +} + fn tier_hint_for_type(tier_type: TierType) -> Option<&'static str> { match tier_type { TierType::RustFS => Some("rustfs"), @@ -1840,10 +2106,10 @@ impl TierConfigMgr { } pub async fn get_driver<'a>(&'a mut self, tier_name: &str) -> std::result::Result<&'a WarmBackendImpl, AdminError> { - if registered_tier_driver_runtime(self).is_some_and(|runtime| lock_unpoisoned(&runtime).draining.contains_key(tier_name)) + if registered_tier_driver_runtime(self).is_some_and(|runtime| runtime_blocks_tier(&lock_unpoisoned(&runtime), tier_name)) { let mut err = ERR_TIER_INVALID_CONFIG.clone(); - err.message = "Remote tier configuration is being replaced".to_string(); + err.message = TIER_MUTATION_BLOCK_MESSAGE.to_string(); return Err(err); } // Return cached driver if present @@ -1871,9 +2137,9 @@ impl TierConfigMgr { let manager = handle.read().await; let runtime = tier_driver_runtime(handle, &manager); let runtime_guard = lock_unpoisoned(&runtime); - if runtime_guard.draining.contains_key(tier_name) { + if runtime_blocks_tier(&runtime_guard, tier_name) { let mut err = ERR_TIER_INVALID_CONFIG.clone(); - err.message = "Remote tier configuration is being replaced".to_string(); + err.message = TIER_MUTATION_BLOCK_MESSAGE.to_string(); return Err(err); } if let Some(generation) = runtime_guard @@ -1891,9 +2157,9 @@ impl TierConfigMgr { let (config, config_fingerprint) = { let mut manager = handle.write().await; let runtime = tier_driver_runtime(handle, &manager); - if lock_unpoisoned(&runtime).draining.contains_key(tier_name) { + if runtime_blocks_tier(&lock_unpoisoned(&runtime), tier_name) { let mut err = ERR_TIER_INVALID_CONFIG.clone(); - err.message = "Remote tier configuration is being replaced".to_string(); + err.message = TIER_MUTATION_BLOCK_MESSAGE.to_string(); return Err(err); } if let Some(driver) = manager.driver_cache.remove(tier_name) { @@ -1926,9 +2192,9 @@ impl TierConfigMgr { let mut manager = handle.write().await; let runtime = tier_driver_runtime(handle, &manager); let runtime_guard = lock_unpoisoned(&runtime); - if runtime_guard.draining.contains_key(tier_name) { + if runtime_blocks_tier(&runtime_guard, tier_name) { let mut err = ERR_TIER_INVALID_CONFIG.clone(); - err.message = "Remote tier configuration is being replaced".to_string(); + err.message = TIER_MUTATION_BLOCK_MESSAGE.to_string(); return Err(err); } if let Some(generation) = runtime_guard @@ -1977,6 +2243,39 @@ impl TierConfigMgr { update_lock.lock_owned().await } + async fn acquire_tier_config_write_lock( + api: Arc, + ) -> std::result::Result + where + S: NamespaceLocking + 'static, + { + let config_file = tier_config_lock_path(); + let ns_lock = api + .new_ns_lock(RUSTFS_META_BUCKET, &config_file) + .await + .map_err(|err| TierConfigUpdateError::Load(io::Error::other(err)))?; + ns_lock + .get_write_lock(get_lock_acquire_timeout()) + .await + .map_err(|err| TierConfigUpdateError::Load(io::Error::other(err))) + } + + async fn update_candidate_with_config_lock( + handle: &Arc>, + api: Arc, + mutation: TierCandidateMutation, + ) -> std::result::Result<(), TierConfigUpdateError> + where + S: TierReferenceProofStore + NamespaceLocking + 'static, + { + let config_lock = Self::acquire_tier_config_write_lock(api.clone()).await?; + let update = Self::admin_update_lock(handle).await; + let (candidate, version) = load_tier_config_for_update(api.clone()) + .await + .map_err(TierConfigUpdateError::Load)?; + Self::update_candidate_owned(handle, api, candidate, version, mutation, update, Some(config_lock)).await + } + #[cfg(test)] pub(crate) async fn publish_candidate( handle: &Arc>, @@ -2033,7 +2332,7 @@ impl TierConfigMgr { err })?; let mut runtime_guard = lock_unpoisoned(&runtime); - if changed.iter().any(|tier_name| runtime_guard.draining.contains_key(tier_name)) { + if changed.iter().any(|tier_name| runtime_blocks_tier(&runtime_guard, tier_name)) { let mut err = ERR_TIER_INVALID_CONFIG.clone(); err.message = "Remote tier configuration is already being replaced".to_string(); return Err(err); @@ -2205,14 +2504,27 @@ impl TierConfigMgr { version: Option, mutation: TierCandidateMutation, update: tokio::sync::OwnedMutexGuard<()>, + config_lock: Option, ) -> std::result::Result<(), TierConfigUpdateError> where - S: EcstoreObjectIO + 'static, + S: TierReferenceProofStore + 'static, { let handle = handle.clone(); tokio::spawn(async move { match AssertUnwindSafe(async move { + let _config_lock = config_lock; let _update = update; + let mutation_kind = mutation.intent_kind(); + let explicit_tier_name = mutation.explicit_tier_name().map(str::to_string); + let current_for_targets = TierConfigMgr { + driver_cache: HashMap::new(), + tiers: candidate + .tiers + .iter() + .map(|(tier_name, config)| (tier_name.clone(), config.clone_with_credentials())) + .collect(), + last_refreshed_at: candidate.last_refreshed_at, + }; let transition = { let mut manager = handle.write().await; let target_tiers = mutation.target_tiers(&manager, &candidate); @@ -2226,6 +2538,14 @@ impl TierConfigMgr { apply_tier_candidate_mutation(mutation, &mut candidate, Instant::now() + TIER_REMOTE_VALIDATION_TIMEOUT) .await .map_err(TierConfigUpdateError::Mutation)?; + let proof_targets = + tier_mutation_proof_targets(mutation_kind, explicit_tier_name.as_deref(), ¤t_for_targets, &candidate); + let affected_targets = + build_tier_mutation_affected_targets(mutation_kind, proof_targets, ¤t_for_targets, &candidate) + .map_err(TierConfigUpdateError::Publish)?; + ensure_no_authoritative_tier_object_references(api.clone(), &affected_targets) + .await + .map_err(TierConfigUpdateError::Publish)?; candidate .save_tiering_config_if_current(api, version.as_deref()) .await @@ -2255,25 +2575,113 @@ impl TierConfigMgr { } pub async fn reload_handle(handle: &Arc>, api: Arc) -> io::Result<()> { + let prepared_intents = Self::load_prepared_mutation_intents(api.clone()).await?; let update = Self::admin_update_lock(handle).await; let candidate = load_tier_config(api).await?; Self::publish_candidate_owned(handle, candidate, None, update) + .await + .map_err(io::Error::other)?; + Self::reconcile_prepared_mutation_intents(handle, &prepared_intents) .await .map_err(io::Error::other) } + async fn load_prepared_mutation_intents(api: Arc) -> io::Result> { + let mut marker = None; + let mut prepared = Vec::new(); + loop { + let scan = list_tier_mutation_intent_records(api.clone(), TIER_MUTATION_INTENT_RECOVERY_SCAN_LIMIT, marker) + .await + .map_err(io::Error::other)?; + if scan.failed != 0 { + return Err(io::Error::other(format!( + "failed to scan {} tier mutation intent record(s) during recovery", + scan.failed + ))); + } + prepared.extend( + scan.intents + .into_iter() + .filter(|intent| intent.state == TierMutationIntentState::Prepared), + ); + if !scan.truncated { + return Ok(prepared); + } + let next_marker = scan.next_marker.ok_or_else(|| { + io::Error::other("tier mutation intent recovery scan was truncated without a continuation marker") + })?; + marker = Some(next_marker); + } + } + + async fn reconcile_prepared_mutation_intents( + handle: &Arc>, + intents: &[TierMutationIntent], + ) -> std::result::Result<(), AdminError> { + let mut prepared_mutation_blocks = HashMap::new(); + for intent in intents { + Self::collect_prepared_mutation_intent_block(&mut prepared_mutation_blocks, intent)?; + } + let manager = handle.read().await; + let runtime = tier_driver_runtime(handle, &manager); + let mut runtime = lock_unpoisoned(&runtime); + runtime.prepared_mutation_blocks = prepared_mutation_blocks; + Ok(()) + } + + pub(crate) async fn apply_prepared_mutation_intent_block( + handle: &Arc>, + intent: &TierMutationIntent, + ) -> std::result::Result<(), AdminError> { + let manager = handle.read().await; + let runtime = tier_driver_runtime(handle, &manager); + let mut runtime = lock_unpoisoned(&runtime); + let mut prepared_mutation_blocks = runtime.prepared_mutation_blocks.clone(); + Self::collect_prepared_mutation_intent_block(&mut prepared_mutation_blocks, intent)?; + runtime.prepared_mutation_blocks = prepared_mutation_blocks; + Ok(()) + } + + pub(crate) async fn clear_prepared_mutation_intent_block(handle: &Arc>, mutation_id: uuid::Uuid) { + let manager = handle.read().await; + let Some(runtime) = registered_tier_driver_runtime(&manager) else { + return; + }; + lock_unpoisoned(&runtime) + .prepared_mutation_blocks + .retain(|_, blocked_mutation_id| *blocked_mutation_id != mutation_id); + } + + fn collect_prepared_mutation_intent_block( + prepared_mutation_blocks: &mut HashMap, + intent: &TierMutationIntent, + ) -> std::result::Result<(), AdminError> { + if intent.state != TierMutationIntentState::Prepared { + return Ok(()); + } + for target in &intent.affected_targets { + match prepared_mutation_blocks.entry(target.tier_name.clone()) { + Entry::Vacant(entry) => { + entry.insert(intent.mutation_id); + } + Entry::Occupied(entry) if *entry.get() == intent.mutation_id => {} + Entry::Occupied(_) => { + let mut err = ERR_TIER_BACKEND_IN_USE.clone(); + err.message = format!("Remote tier {} already has another prepared mutation", target.tier_name); + return Err(err); + } + } + } + Ok(()) + } + pub async fn add_and_save( handle: &Arc>, api: Arc, tier_config: TierConfig, force: bool, ) -> std::result::Result<(), TierConfigUpdateError> { - let update = Self::admin_update_lock(handle).await; - let (candidate, version) = load_tier_config_for_update(api.clone()) - .await - .map_err(TierConfigUpdateError::Load)?; - Self::update_candidate_owned(handle, api, candidate, version, TierCandidateMutation::Add(tier_config, force), update) - .await + Self::update_candidate_with_config_lock(handle, api, TierCandidateMutation::Add(tier_config, force)).await } pub async fn edit_and_save( @@ -2282,19 +2690,8 @@ impl TierConfigMgr { tier_name: &str, credentials: TierCreds, ) -> std::result::Result<(), TierConfigUpdateError> { - let update = Self::admin_update_lock(handle).await; - let (candidate, version) = load_tier_config_for_update(api.clone()) + Self::update_candidate_with_config_lock(handle, api, TierCandidateMutation::Edit(tier_name.to_string(), credentials)) .await - .map_err(TierConfigUpdateError::Load)?; - Self::update_candidate_owned( - handle, - api, - candidate, - version, - TierCandidateMutation::Edit(tier_name.to_string(), credentials), - update, - ) - .await } pub async fn remove_and_save( @@ -2303,7 +2700,7 @@ impl TierConfigMgr { tier_name: &str, force: bool, ) -> std::result::Result<(), TierConfigUpdateError> { - Self::remove_and_save_with(handle, api, tier_name, force).await + Self::update_candidate_with_config_lock(handle, api, TierCandidateMutation::Remove(tier_name.to_string(), force)).await } async fn remove_and_save_with( @@ -2313,7 +2710,7 @@ impl TierConfigMgr { force: bool, ) -> std::result::Result<(), TierConfigUpdateError> where - S: EcstoreObjectIO + 'static, + S: TierReferenceProofStore + 'static, { let update = Self::admin_update_lock(handle).await; let (candidate, version) = load_tier_config_for_update(api.clone()) @@ -2326,6 +2723,7 @@ impl TierConfigMgr { version, TierCandidateMutation::Remove(tier_name.to_string(), force), update, + None, ) .await } @@ -2335,7 +2733,7 @@ impl TierConfigMgr { api: Arc, force: bool, ) -> std::result::Result<(), TierConfigUpdateError> { - Self::clear_and_save_with(handle, api, force).await + Self::update_candidate_with_config_lock(handle, api, TierCandidateMutation::Clear(force)).await } async fn clear_and_save_with( @@ -2344,13 +2742,13 @@ impl TierConfigMgr { force: bool, ) -> std::result::Result<(), TierConfigUpdateError> where - S: EcstoreObjectIO + 'static, + S: TierReferenceProofStore + 'static, { let update = Self::admin_update_lock(handle).await; let (candidate, version) = load_tier_config_for_update(api.clone()) .await .map_err(TierConfigUpdateError::Load)?; - Self::update_candidate_owned(handle, api, candidate, version, TierCandidateMutation::Clear(force), update).await + Self::update_candidate_owned(handle, api, candidate, version, TierCandidateMutation::Clear(force), update, None).await } pub async fn verify_without_manager_lock(handle: &Arc>, tier_name: &str) -> std::result::Result<(), io::Error> { @@ -2467,7 +2865,7 @@ impl TierConfigMgr { return Ok(()); }; let runtime = lock_unpoisoned(&runtime); - let blocked = runtime.draining.contains_key(tier_name) + let blocked = runtime_blocks_tier(&runtime, tier_name) || runtime .generations .get(tier_name) @@ -2555,7 +2953,7 @@ impl TierConfigMgr { } fn apply_reloaded_tiers(&mut self, tiers: HashMap) -> std::result::Result<(), AdminError> { - if registered_tier_driver_runtime(self).is_some_and(|runtime| !lock_unpoisoned(&runtime).draining.is_empty()) { + if registered_tier_driver_runtime(self).is_some_and(|runtime| runtime_has_mutation_block(&lock_unpoisoned(&runtime))) { return Err(ERR_TIER_BACKEND_IN_USE.clone()); } let changed_or_removed = self @@ -2585,7 +2983,7 @@ impl TierConfigMgr { Err(err) => { if let Some(runtime) = registered_tier_driver_runtime(self) { let runtime = lock_unpoisoned(&runtime); - if !runtime.draining.is_empty() + if runtime_has_mutation_block(&runtime) || runtime .generations .values() @@ -2604,7 +3002,7 @@ impl TierConfigMgr { } pub async fn clear_tier(&mut self, force: bool) -> std::result::Result<(), AdminError> { - if registered_tier_driver_runtime(self).is_some_and(|runtime| !lock_unpoisoned(&runtime).draining.is_empty()) { + if registered_tier_driver_runtime(self).is_some_and(|runtime| runtime_has_mutation_block(&lock_unpoisoned(&runtime))) { return Err(ERR_TIER_BACKEND_IN_USE.clone()); } self.ensure_generations_are_idle(self.tiers.keys())?; @@ -3484,6 +3882,238 @@ mod tests { } } + #[derive(Debug)] + struct LockingTierConfigStore { + locks: Mutex>, + put_after_lock: AtomicBool, + lock_manager: Arc, + } + + impl LockingTierConfigStore { + fn new() -> Self { + Self { + locks: Mutex::new(Vec::new()), + put_after_lock: AtomicBool::new(false), + lock_manager: Arc::new(rustfs_lock::GlobalLockManager::new()), + } + } + + fn lock_events(&self) -> Vec<(String, String)> { + self.locks + .lock() + .expect("tier config lock recorder should not poison") + .clone() + } + } + + #[async_trait::async_trait] + impl ObjectIO for LockingTierConfigStore { + type Error = Error; + type RangeSpec = HTTPRangeSpec; + type HeaderMap = HeaderMap; + type ObjectOptions = ObjectOptions; + type ObjectInfo = ObjectInfo; + type GetObjectReader = GetObjectReader; + type PutObjectReader = PutObjReader; + + async fn get_object_reader( + &self, + bucket: &str, + object: &str, + _range: Option, + _h: HeaderMap, + _opts: &ObjectOptions, + ) -> Result { + Err(Error::ObjectNotFound(bucket.to_string(), object.to_string())) + } + + async fn put_object( + &self, + bucket: &str, + object: &str, + _data: &mut PutObjReader, + opts: &ObjectOptions, + ) -> Result { + let expected_lock = (RUSTFS_META_BUCKET.to_string(), tier_config_lock_path()); + let saw_expected_lock = self + .locks + .lock() + .expect("tier config lock recorder should not poison") + .contains(&expected_lock); + self.put_after_lock.store(saw_expected_lock, Ordering::SeqCst); + assert!(opts.max_parity, "tier config save must retain max parity"); + assert!( + opts.http_preconditions + .as_ref() + .and_then(|preconditions| preconditions.if_none_match.as_deref()) + == Some("*"), + "initial tier config save must retain if-none-match protection" + ); + Ok(ObjectInfo { + bucket: bucket.to_string(), + name: object.to_string(), + etag: Some("saved-etag".to_string()), + ..Default::default() + }) + } + } + + #[async_trait::async_trait] + impl BucketOperations for LockingTierConfigStore { + type Error = Error; + + async fn make_bucket( + &self, + _bucket: &str, + _opts: &crate::storage_api_contracts::bucket::MakeBucketOptions, + ) -> Result<()> { + Err(Error::NotImplemented) + } + + async fn get_bucket_info( + &self, + _bucket: &str, + _opts: &crate::storage_api_contracts::bucket::BucketOptions, + ) -> Result { + Err(Error::NotImplemented) + } + + async fn list_bucket( + &self, + _opts: &crate::storage_api_contracts::bucket::BucketOptions, + ) -> Result> { + Ok(Vec::new()) + } + + async fn delete_bucket( + &self, + _bucket: &str, + _opts: &crate::storage_api_contracts::bucket::DeleteBucketOptions, + ) -> Result<()> { + Err(Error::NotImplemented) + } + } + + #[async_trait::async_trait] + impl ListOperations for LockingTierConfigStore { + type Error = Error; + type ListObjectsV2Info = StorageListObjectsV2Info; + type ListObjectVersionsInfo = StorageListObjectVersionsInfo; + type ObjectInfoOrErr = StorageObjectInfoOrErr; + type WalkOptions = TierReferenceProofWalkOptions; + type WalkCancellation = tokio_util::sync::CancellationToken; + type WalkResultSender = tokio::sync::mpsc::Sender>; + + async fn list_objects_v2( + self: Arc, + _bucket: &str, + _prefix: &str, + _continuation_token: Option, + _delimiter: Option, + _max_keys: i32, + _fetch_owner: bool, + _start_after: Option, + _incl_deleted: bool, + ) -> Result { + Ok(StorageListObjectsV2Info::default()) + } + + async fn list_object_versions( + self: Arc, + _bucket: &str, + _prefix: &str, + _marker: Option, + _version_marker: Option, + _delimiter: Option, + _max_keys: i32, + ) -> Result { + Ok(StorageListObjectVersionsInfo::default()) + } + + async fn walk( + self: Arc, + _rx: Self::WalkCancellation, + _bucket: &str, + _prefix: &str, + _result: Self::WalkResultSender, + _opts: Self::WalkOptions, + ) -> Result<()> { + Err(Error::NotImplemented) + } + } + + #[async_trait::async_trait] + impl NamespaceLocking for LockingTierConfigStore { + type Error = Error; + type NamespaceLock = rustfs_lock::NamespaceLockWrapper; + + async fn new_ns_lock(&self, bucket: &str, object: &str) -> Result { + self.locks + .lock() + .expect("tier config lock recorder should not poison") + .push((bucket.to_string(), object.to_string())); + Ok(rustfs_lock::NamespaceLockWrapper::new( + rustfs_lock::NamespaceLock::with_local_manager("tier-config-test".to_string(), self.lock_manager.clone()), + rustfs_lock::ObjectKey { + bucket: Arc::from(bucket), + object: Arc::from(object), + version: None, + }, + "tier-config-test-owner".to_string(), + )) + } + } + + #[tokio::test] + async fn tier_config_update_path_acquires_meta_namespace_sidecar_lock_before_save() { + let manager = TierConfigMgr::new(); + let store = Arc::new(LockingTierConfigStore::new()); + + TierConfigMgr::update_candidate_with_config_lock(&manager, store.clone(), TierCandidateMutation::Clear(true)) + .await + .expect("empty clear should still acquire and save through the coordinator lock"); + + assert_eq!(store.lock_events(), vec![(RUSTFS_META_BUCKET.to_string(), tier_config_lock_path())]); + assert!( + store.put_after_lock.load(Ordering::SeqCst), + "tier config save must happen after the coordinator namespace lock is acquired" + ); + } + + #[test] + fn tier_mutation_targets_prove_remove_and_rebind_authoritative_identity() { + let mut current = empty_mgr(); + current.tiers.insert("COLD-A".to_string(), build_rustfs_tier("COLD-A")); + + let mut rebound = empty_mgr(); + let mut replacement = build_rustfs_tier("COLD-A"); + replacement.rustfs.as_mut().expect("replacement payload should exist").prefix = "new-prefix".to_string(); + rebound.tiers.insert("COLD-A".to_string(), replacement); + + let edit_targets = TierCandidateMutation::Edit("COLD-A".to_string(), TierCreds::default()) + .affected_targets(¤t, &rebound) + .expect("rebind target proof should build"); + assert_eq!(edit_targets.len(), 1); + assert_eq!(edit_targets[0].tier_name, "COLD-A"); + assert!(edit_targets[0].old_backend_identity.is_some()); + assert!(edit_targets[0].new_backend_identity.is_some()); + assert_ne!(edit_targets[0].old_backend_identity, edit_targets[0].new_backend_identity); + + let removed = empty_mgr(); + let remove_targets = TierCandidateMutation::Remove("COLD-A".to_string(), true) + .affected_targets(¤t, &removed) + .expect("remove target proof should build"); + assert_eq!(remove_targets.len(), 1); + assert_eq!(remove_targets[0].tier_name, "COLD-A"); + assert!(remove_targets[0].old_backend_identity.is_some()); + assert!(remove_targets[0].new_backend_identity.is_none()); + + let noop_targets = TierCandidateMutation::Remove("MISSING".to_string(), true) + .affected_targets(¤t, ¤t) + .expect("idempotent remove of a missing tier should not require an identity proof"); + assert!(noop_targets.is_empty()); + } + /// A fully offline `WarmBackend` used to exercise the driver-facing /// branches of `remove`/`verify` without touching a remote tier. struct MockWarmBackend { @@ -4330,6 +4960,110 @@ mod tests { .expect("test driver generation should install"); } + fn prepared_remove_intent(tier_name: &str, mutation_id: uuid::Uuid) -> TierMutationIntent { + TierMutationIntent { + mutation_id, + revision: 1, + kind: TierMutationIntentKind::Remove, + state: TierMutationIntentState::Prepared, + old_config_etag: Some("etag-old".to_string()), + committed_config_etag: None, + candidate_digest: [1; 32], + affected_targets: vec![TierMutationIntentTarget { + tier_name: tier_name.to_string(), + old_backend_identity: Some([2; 32]), + new_backend_identity: None, + }], + expires_at_unix_nanos: 1, + } + } + + #[tokio::test] + async fn prepared_mutation_recovery_blocks_new_operation_lease() { + let manager = TierConfigMgr::new(); + manager + .write() + .await + .tiers + .insert("COLD-A".to_string(), build_rustfs_tier("COLD-A")); + let mutation_id = uuid::Uuid::from_u128(1); + TierConfigMgr::reconcile_prepared_mutation_intents(&manager, &[prepared_remove_intent("COLD-A", mutation_id)]) + .await + .expect("prepared recovery block should apply"); + + let err = match TierConfigMgr::acquire_operation_lease(&manager, "COLD-A").await { + Ok(_) => panic!("prepared mutation must block new operation leases"), + Err(err) => err, + }; + assert_eq!(err.code, ERR_TIER_INVALID_CONFIG.code); + assert_eq!(err.message, TIER_MUTATION_BLOCK_MESSAGE); + let guard = manager.read().await; + let runtime = registered_tier_driver_runtime(&guard).expect("runtime should be registered by recovery block"); + assert_eq!(lock_unpoisoned(&runtime).prepared_mutation_blocks.get("COLD-A"), Some(&mutation_id)); + } + + #[tokio::test] + async fn prepared_mutation_recovery_blocks_admin_publish() { + let manager = TierConfigMgr::new(); + manager + .write() + .await + .tiers + .insert("COLD-A".to_string(), build_rustfs_tier("COLD-A")); + TierConfigMgr::reconcile_prepared_mutation_intents( + &manager, + &[prepared_remove_intent("COLD-A", uuid::Uuid::from_u128(1))], + ) + .await + .expect("prepared recovery block should apply"); + let candidate = empty_mgr(); + + let err = TierConfigMgr::publish_candidate(&manager, candidate, None) + .await + .expect_err("prepared mutation must block conflicting admin publishes"); + assert_eq!(err.code, ERR_TIER_INVALID_CONFIG.code); + assert!(manager.read().await.tiers.contains_key("COLD-A")); + } + + #[tokio::test] + async fn prepared_mutation_recovery_is_idempotent_and_rejects_conflicts() { + let manager = TierConfigMgr::new(); + let first = prepared_remove_intent("COLD-A", uuid::Uuid::from_u128(1)); + TierConfigMgr::reconcile_prepared_mutation_intents(&manager, &[first.clone(), first]) + .await + .expect("same prepared mutation should be idempotent"); + + let first = prepared_remove_intent("COLD-A", uuid::Uuid::from_u128(1)); + let second = prepared_remove_intent("COLD-A", uuid::Uuid::from_u128(2)); + let err = TierConfigMgr::reconcile_prepared_mutation_intents(&manager, &[first, second]) + .await + .expect_err("a second prepared mutation for the same tier must fail closed"); + assert_eq!(err.code, ERR_TIER_BACKEND_IN_USE.code); + } + + #[tokio::test] + async fn prepared_mutation_recovery_clears_resolved_intents() { + let manager = TierConfigMgr::new(); + manager + .write() + .await + .tiers + .insert("COLD-A".to_string(), build_rustfs_tier("COLD-A")); + TierConfigMgr::reconcile_prepared_mutation_intents( + &manager, + &[prepared_remove_intent("COLD-A", uuid::Uuid::from_u128(1))], + ) + .await + .expect("prepared recovery block should apply"); + TierConfigMgr::reconcile_prepared_mutation_intents(&manager, &[]) + .await + .expect("resolved recovery scan should clear stale blocks"); + + let guard = manager.read().await; + let runtime = registered_tier_driver_runtime(&guard).expect("runtime should stay registered"); + assert!(lock_unpoisoned(&runtime).prepared_mutation_blocks.is_empty()); + } + #[tokio::test] async fn registry_steady_state_lookup_does_not_scan_misses() { let manager = TierConfigMgr::new(); @@ -5324,11 +6058,36 @@ mod tests { assert!(lock_unpoisoned(&runtime).generations.get("COLD-A").is_none()); } - #[derive(Debug, Default)] + #[derive(Debug)] struct CasConfigStore { state: tokio::sync::Mutex, String)>>, next_etag: AtomicUsize, fail_put: AtomicBool, + lock_manager: Arc, + lock_requests: Mutex>, + listed_versions: Mutex>, + } + + impl Default for CasConfigStore { + fn default() -> Self { + Self { + state: tokio::sync::Mutex::new(None), + next_etag: AtomicUsize::new(0), + fail_put: AtomicBool::new(false), + lock_manager: Arc::new(rustfs_lock::GlobalLockManager::new()), + lock_requests: Mutex::new(Vec::new()), + listed_versions: Mutex::new(Vec::new()), + } + } + } + + impl CasConfigStore { + fn add_listed_version(&self, object: ObjectInfo) { + self.listed_versions + .lock() + .expect("tier reference fixture should not poison") + .push(object); + } } #[async_trait::async_trait] @@ -5408,6 +6167,195 @@ mod tests { } } + #[async_trait::async_trait] + impl BucketOperations for CasConfigStore { + type Error = Error; + + async fn make_bucket( + &self, + _bucket: &str, + _opts: &crate::storage_api_contracts::bucket::MakeBucketOptions, + ) -> Result<()> { + Err(Error::NotImplemented) + } + + async fn get_bucket_info( + &self, + _bucket: &str, + _opts: &crate::storage_api_contracts::bucket::BucketOptions, + ) -> Result { + Err(Error::NotImplemented) + } + + async fn list_bucket( + &self, + _opts: &crate::storage_api_contracts::bucket::BucketOptions, + ) -> Result> { + let mut seen = HashSet::new(); + let mut buckets = Vec::new(); + for object in self + .listed_versions + .lock() + .expect("tier reference fixture should not poison") + .iter() + { + if seen.insert(object.bucket.clone()) { + buckets.push(crate::storage_api_contracts::bucket::BucketInfo { + name: object.bucket.clone(), + ..Default::default() + }); + } + } + Ok(buckets) + } + + async fn delete_bucket( + &self, + _bucket: &str, + _opts: &crate::storage_api_contracts::bucket::DeleteBucketOptions, + ) -> Result<()> { + Err(Error::NotImplemented) + } + } + + #[async_trait::async_trait] + impl ListOperations for CasConfigStore { + type Error = Error; + type ListObjectsV2Info = StorageListObjectsV2Info; + type ListObjectVersionsInfo = StorageListObjectVersionsInfo; + type ObjectInfoOrErr = StorageObjectInfoOrErr; + type WalkOptions = TierReferenceProofWalkOptions; + type WalkCancellation = tokio_util::sync::CancellationToken; + type WalkResultSender = tokio::sync::mpsc::Sender>; + + async fn list_objects_v2( + self: Arc, + _bucket: &str, + _prefix: &str, + _continuation_token: Option, + _delimiter: Option, + _max_keys: i32, + _fetch_owner: bool, + _start_after: Option, + _incl_deleted: bool, + ) -> Result { + Ok(StorageListObjectsV2Info::default()) + } + + async fn list_object_versions( + self: Arc, + bucket: &str, + prefix: &str, + marker: Option, + version_marker: Option, + _delimiter: Option, + max_keys: i32, + ) -> Result { + let mut objects: Vec<_> = self + .listed_versions + .lock() + .expect("tier reference fixture should not poison") + .iter() + .filter(|object| object.bucket == bucket && object.name.starts_with(prefix)) + .cloned() + .collect(); + objects.sort_by(|left, right| tier_test_object_marker(left).cmp(&tier_test_object_marker(right))); + if marker.is_some() || version_marker.is_some() { + let marker = (marker.unwrap_or_default(), version_marker.unwrap_or_default()); + objects.retain(|object| tier_test_object_marker(object) > marker); + } + let limit = match usize::try_from(max_keys) { + Ok(limit) => limit, + Err(_) => 0, + }; + let is_truncated = objects.len() > limit; + if is_truncated { + objects.truncate(limit); + } + let (next_marker, next_version_idmarker) = if is_truncated { + objects + .last() + .map(|object| (Some(object.name.clone()), object.version_id.map(|version| version.to_string()))) + .unwrap_or((None, None)) + } else { + (None, None) + }; + Ok(StorageListObjectVersionsInfo { + is_truncated, + next_marker, + next_version_idmarker, + objects, + prefixes: Vec::new(), + }) + } + + async fn walk( + self: Arc, + _rx: Self::WalkCancellation, + _bucket: &str, + _prefix: &str, + _result: Self::WalkResultSender, + _opts: Self::WalkOptions, + ) -> Result<()> { + Err(Error::NotImplemented) + } + } + + fn tier_test_object_marker(object: &ObjectInfo) -> (String, String) { + ( + object.name.clone(), + object.version_id.map(|version| version.to_string()).unwrap_or_default(), + ) + } + + fn transitioned_tier_object( + bucket: &str, + object: &str, + tier_name: &str, + backend_identity: Option, + ) -> ObjectInfo { + let mut user_defined = HashMap::new(); + if let Some(identity) = backend_identity { + rustfs_utils::http::metadata_compat::insert_str( + &mut user_defined, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID, + rustfs_utils::crypto::hex(identity), + ); + } + ObjectInfo { + bucket: bucket.to_string(), + name: object.to_string(), + user_defined: Arc::new(user_defined), + transitioned_object: crate::storage_api_contracts::lifecycle::TransitionedObject { + name: object.to_string(), + tier: tier_name.to_string(), + status: rustfs_filemeta::TRANSITION_COMPLETE.to_string(), + ..Default::default() + }, + ..Default::default() + } + } + + #[async_trait::async_trait] + impl NamespaceLocking for CasConfigStore { + type Error = Error; + type NamespaceLock = rustfs_lock::NamespaceLockWrapper; + + async fn new_ns_lock(&self, bucket: &str, object: &str) -> Result { + self.lock_requests + .lock() + .expect("tier config lock request log should not poison") + .push((bucket.to_string(), object.to_string())); + let lock = + rustfs_lock::NamespaceLock::with_local_manager("tier-config-update-test".to_string(), self.lock_manager.clone()); + Ok(rustfs_lock::NamespaceLockWrapper::new( + lock, + rustfs_lock::ObjectKey::new(bucket.to_string(), object.to_string()), + "tier-config-update-test-owner".to_string(), + )) + } + } + #[tokio::test] async fn remove_and_clear_full_update_paths_preserve_force() { let remove_store = Arc::new(CasConfigStore::default()); @@ -5461,6 +6409,151 @@ mod tests { ); } + #[tokio::test] + async fn zero_reference_proof_blocks_remove_before_config_save() { + let store = Arc::new(CasConfigStore::default()); + let tier = build_rustfs_tier("COLD-A"); + let identity = tier_backend_identity(&tier).expect("test tier identity should encode"); + let mut persisted = empty_mgr(); + persisted.tiers.insert("COLD-A".to_string(), tier.clone_with_credentials()); + persisted + .save_tiering_config_if_current(store.clone(), None) + .await + .expect("reference proof fixture should persist"); + store.add_listed_version(transitioned_tier_object("photos", "2026/a.jpg", "COLD-A", Some(identity))); + + let manager = TierConfigMgr::new(); + manager.write().await.tiers.insert("COLD-A".to_string(), tier); + let err = TierConfigMgr::remove_and_save_with(&manager, store.clone(), "COLD-A", true) + .await + .expect_err("authoritative object reference must block tier removal"); + + match err { + TierConfigUpdateError::Publish(err) => { + assert_eq!(err.code, ERR_TIER_BACKEND_IN_USE.code); + assert!(err.message.contains("photos/2026/a.jpg"), "{}", err.message); + } + other => panic!("reference proof failures must be surfaced as publish errors: {other:?}"), + } + assert!(manager.read().await.tiers.contains_key("COLD-A")); + assert!( + load_tier_config_for_update(store) + .await + .expect("blocked config should still reload") + .0 + .tiers + .contains_key("COLD-A"), + "blocked removal must not persist the empty candidate" + ); + } + + #[test] + fn zero_reference_proof_rebind_identity_boundaries_fail_closed() { + let tier = build_rustfs_tier("COLD-A"); + let current_identity = tier_backend_identity(&tier).expect("current identity should encode"); + let replacement_identity = + tier_backend_identity(&build_azure_tier("account-a")).expect("replacement identity should encode"); + let target = TierMutationIntentTarget { + tier_name: "COLD-A".to_string(), + old_backend_identity: Some(current_identity), + new_backend_identity: Some(current_identity), + }; + let object = transitioned_tier_object("photos", "2026/a.jpg", "COLD-A", Some(current_identity)); + assert!(!tier_object_blocks_target_rebind(&object, &target).expect("matching identity should parse")); + + let target = TierMutationIntentTarget { + tier_name: "COLD-A".to_string(), + old_backend_identity: Some(current_identity), + new_backend_identity: Some(replacement_identity), + }; + assert!(tier_object_blocks_target_rebind(&object, &target).expect("mismatched identity should parse")); + + let object_without_identity = transitioned_tier_object("photos", "2026/b.jpg", "COLD-A", None); + assert!( + tier_object_blocks_target_rebind(&object_without_identity, &target) + .expect("missing identity means unknown reference owner") + ); + + let mut pending = transitioned_tier_object("photos", "2026/c.jpg", "COLD-A", Some(current_identity)); + pending.transitioned_object.status = "pending".to_string(); + assert!(!tier_object_blocks_target_rebind(&pending, &target).expect("non-complete status should not block")); + } + + #[tokio::test] + async fn mutation_target_proof_uses_persisted_config_snapshot_not_stale_manager() { + let manager = TierConfigMgr::new(); + let store = Arc::new(CasConfigStore::default()); + let mut candidate = empty_mgr(); + candidate.tiers.insert("COLD-A".to_string(), build_rustfs_tier("COLD-A")); + candidate.tiers.insert("COLD-B".to_string(), build_rustfs_tier("COLD-B")); + let update = TierConfigMgr::admin_update_lock(&manager).await; + + TierConfigMgr::update_candidate_owned( + &manager, + store.clone(), + candidate, + None, + TierCandidateMutation::Remove("COLD-A".to_string(), true), + update, + None, + ) + .await + .expect("authoritative proof should use the loaded config even when the local manager is stale"); + assert!(!manager.read().await.tiers.contains_key("COLD-A")); + assert!(manager.read().await.tiers.contains_key("COLD-B")); + let reloaded = load_tier_config_for_update(store) + .await + .expect("removed config should reload") + .0; + assert!(!reloaded.tiers.contains_key("COLD-A")); + assert!(reloaded.tiers.contains_key("COLD-B")); + } + + #[test] + fn add_target_proof_ignores_unchanged_persisted_tiers_when_manager_is_stale() { + let mut current = empty_mgr(); + current.tiers.insert("COLD-A".to_string(), build_rustfs_tier("COLD-A")); + + let mut candidate = empty_mgr(); + candidate.tiers.insert("COLD-A".to_string(), build_rustfs_tier("COLD-A")); + candidate.tiers.insert("COLD-B".to_string(), build_rustfs_tier("COLD-B")); + + let targets = TierCandidateMutation::Add(build_rustfs_tier("COLD-B"), true) + .affected_targets(¤t, &candidate) + .expect("add proof should ignore unchanged durable tiers"); + assert_eq!(targets.len(), 1); + assert_eq!(targets[0].tier_name, "COLD-B"); + assert!(targets[0].old_backend_identity.is_none()); + assert!(targets[0].new_backend_identity.is_some()); + } + + #[tokio::test] + async fn update_with_config_lock_add_uses_loaded_snapshot_for_target_proof() { + let manager = TierConfigMgr::new(); + let store = Arc::new(CasConfigStore::default()); + let mut persisted = empty_mgr(); + persisted.tiers.insert("COLD-A".to_string(), build_rustfs_tier("COLD-A")); + persisted + .save_tiering_config_if_current(store.clone(), None) + .await + .expect("add fixture should persist"); + + TierConfigMgr::update_candidate_with_config_lock( + &manager, + store.clone(), + TierCandidateMutation::Add(build_rustfs_tier("COLD-B"), true), + ) + .await + .expect("add proof must ignore unchanged durable tiers missing from the stale manager"); + + let reloaded = load_tier_config_for_update(store) + .await + .expect("updated config should reload") + .0; + assert!(reloaded.tiers.contains_key("COLD-A")); + assert!(reloaded.tiers.contains_key("COLD-B")); + } + #[tokio::test] async fn failed_owned_update_restores_generation_and_reports_save_error() { let manager = TierConfigMgr::new(); @@ -5484,6 +6577,7 @@ mod tests { None, TierCandidateMutation::Remove("COLD-A".to_string(), true), update, + None, ) .await .expect_err("save failure must be observable to the admin caller"); @@ -5513,6 +6607,7 @@ mod tests { None, TierCandidateMutation::Remove("COLD-A".to_string(), true), update, + None, ) .await .expect("a later update must succeed after save failure recovery"); @@ -5546,6 +6641,7 @@ mod tests { None, TierCandidateMutation::Remove("COLD-A".to_string(), true), update, + None, ) .await }); @@ -5572,6 +6668,106 @@ mod tests { assert!(store.state.lock().await.is_some(), "detached update must persist its result"); } + #[tokio::test] + async fn config_write_lock_survives_cancelled_wrapped_update_until_detached_publish_finishes() { + let manager = TierConfigMgr::new(); + { + let mut guard = manager.write().await; + install_lease_backend(&mut guard, "COLD-A", LeaseTestBackend::ready("old")); + } + let old = TierConfigMgr::acquire_operation_lease(&manager, "COLD-A") + .await + .expect("old generation lease should be available"); + let store = Arc::new(CasConfigStore::default()); + let mut candidate = empty_mgr(); + candidate.tiers.insert("COLD-A".to_string(), build_rustfs_tier("COLD-A")); + candidate + .save_tiering_config_if_current(store.clone(), None) + .await + .expect("tier config fixture should persist"); + let update_manager = manager.clone(); + let update_store = store.clone(); + let caller = tokio::spawn(async move { + TierConfigMgr::update_candidate_with_config_lock( + &update_manager, + update_store, + TierCandidateMutation::Remove("COLD-A".to_string(), true), + ) + .await + }); + tokio::time::timeout(Duration::from_secs(1), async { + while old.is_current(&manager).await { + tokio::task::yield_now().await; + } + }) + .await + .expect("owned update should revoke before caller cancellation"); + caller.abort(); + + let config_file = tier_config_lock_path(); + let competing_lock = store + .new_ns_lock(RUSTFS_META_BUCKET, &config_file) + .await + .expect("competing tier config lock should be created"); + let err = competing_lock + .get_write_lock(Duration::from_millis(20)) + .await + .expect_err("detached update must keep the config write lock until it finishes"); + assert!(matches!(err, rustfs_lock::LockError::Timeout { .. })); + + drop(old); + tokio::time::timeout(Duration::from_secs(1), async { + while manager.read().await.tiers.contains_key("COLD-A") { + tokio::task::yield_now().await; + } + }) + .await + .expect("detached update must finish after its operation lease drains"); + competing_lock + .get_write_lock(Duration::from_secs(1)) + .await + .expect("tier config lock should release after detached update finishes"); + } + + #[tokio::test] + #[serial_test::serial] + async fn tier_config_update_waits_for_config_namespace_lock_before_loading() { + let manager = TierConfigMgr::new(); + let store = Arc::new(CasConfigStore::default()); + let config_file = tier_config_lock_path(); + let outer_lock = store + .new_ns_lock(RUSTFS_META_BUCKET, &config_file) + .await + .expect("outer tier config lock should be created"); + let _outer_guard = outer_lock + .get_write_lock(Duration::from_secs(1)) + .await + .expect("outer tier config lock should be acquired"); + + let result = temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("1"))], async { + TierConfigMgr::update_candidate_with_config_lock(&manager, store.clone(), TierCandidateMutation::Clear(true)).await + }) + .await; + let TierConfigUpdateError::Load(err) = result.expect_err("config mutation must wait behind the config namespace lock") + else { + panic!("config lock contention should fail before candidate load or mutation"); + }; + let rendered = err.to_string(); + assert!(rendered.to_ascii_lowercase().contains("timeout"), "{rendered}"); + assert!(rendered.contains(&config_file), "{rendered}"); + assert_eq!( + store + .lock_requests + .lock() + .expect("tier config lock request log should not poison") + .as_slice(), + &[ + (RUSTFS_META_BUCKET.to_string(), config_file.clone()), + (RUSTFS_META_BUCKET.to_string(), config_file), + ] + ); + } + #[tokio::test] async fn panicked_owned_update_restores_generation_and_reports_error() { let manager = TierConfigMgr::new(); @@ -5594,6 +6790,7 @@ mod tests { None, TierCandidateMutation::Remove("COLD-A".to_string(), false), update, + None, ) .await .expect_err("mutation panic must be observable to the admin caller"); diff --git a/crates/ecstore/src/services/tier/tier_mutation_intent.rs b/crates/ecstore/src/services/tier/tier_mutation_intent.rs index dec3260e0..e842c4650 100644 --- a/crates/ecstore/src/services/tier/tier_mutation_intent.rs +++ b/crates/ecstore/src/services/tier/tier_mutation_intent.rs @@ -27,7 +27,7 @@ use crate::storage_api_contracts::{list::ListOperations as _, object::HTTPPrecon 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 MAX_TIER_MUTATION_INTENT_SIZE: usize = rustfs_protos::TIER_MUTATION_RPC_MAX_PREPARE_PAYLOAD_SIZE; pub(crate) const TIER_MUTATION_INTENT_RECORD_PREFIX: &str = "tier/mutation-intents/records"; pub(crate) type TierMutationDigest = [u8; 32]; @@ -346,6 +346,28 @@ pub(crate) async fn save_tier_mutation_intent_record(api: Arc, intent: com::save_config(api, &object, data).await } +pub(crate) async fn save_tier_mutation_intent_record_if_absent( + 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_with_opts( + api, + &object, + data, + &ObjectOptions { + max_parity: true, + http_preconditions: Some(HTTPPreconditions { + if_none_match: Some("*".to_string()), + ..Default::default() + }), + ..Default::default() + }, + ) + .await +} + pub(crate) async fn load_tier_mutation_intent_record(api: Arc, mutation_id: Uuid) -> EcstoreResult { let (intent, _) = load_tier_mutation_intent_record_with_etag(api, mutation_id).await?; Ok(intent) diff --git a/crates/ecstore/src/services/tier/tier_mutation_peer.rs b/crates/ecstore/src/services/tier/tier_mutation_peer.rs new file mode 100644 index 000000000..5ccfd043b --- /dev/null +++ b/crates/ecstore/src/services/tier/tier_mutation_peer.rs @@ -0,0 +1,278 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::sync::Arc; + +use rustfs_protos::{TIER_MUTATION_RPC_PROTOCOL_VERSION, TierMutationRpcPhase}; +use uuid::Uuid; + +use super::tier::TierConfigMgr; +use super::tier_mutation_intent::{ + MAX_TIER_MUTATION_INTENT_SIZE, TierMutationIntent, TierMutationIntentState, advance_tier_mutation_intent_record_idempotent, + load_tier_mutation_intent_record, save_tier_mutation_intent_record_if_absent, +}; +use crate::client::admin_handler_utils::AdminError; +use crate::error::{Error, StorageError}; +use crate::store::ECStore; + +pub const MAX_TIER_MUTATION_PEER_COMMIT_ETAG_SIZE: usize = rustfs_protos::TIER_MUTATION_RPC_MAX_COMMIT_PAYLOAD_SIZE; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum TierMutationPeerState { + Prepared, + Committed, + Aborted, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct TierMutationPeerOutcome { + pub state: TierMutationPeerState, + pub applied: bool, +} + +#[derive(Debug, thiserror::Error)] +pub enum TierMutationPeerError { + #[error("unsupported tier mutation peer protocol version: {0}")] + UnsupportedProtocolVersion(u32), + #[error("tier mutation peer mutation_id is nil")] + NilMutationId, + #[error("tier mutation peer payload is too large: {len}/{max}")] + PayloadTooLarge { len: usize, max: usize }, + #[error("tier mutation peer payload is invalid: {0}")] + InvalidPayload(String), + #[error("tier mutation peer intent conflicts with existing record")] + ConflictingIntent, + #[error("tier mutation peer runtime error: {0}")] + Runtime(#[source] AdminError), + #[error("tier mutation peer store error: {0}")] + Store(#[source] StorageError), +} + +impl From for TierMutationPeerError { + fn from(error: Error) -> Self { + Self::Store(error) + } +} + +pub type TierMutationPeerResult = std::result::Result; + +pub async fn handle_tier_mutation_peer_request( + api: Arc, + protocol_version: u32, + phase: TierMutationRpcPhase, + mutation_id: Uuid, + canonical_payload: &[u8], +) -> TierMutationPeerResult { + validate_peer_request_envelope(protocol_version, mutation_id, canonical_payload)?; + match phase { + TierMutationRpcPhase::Prepare => handle_prepare(api, mutation_id, canonical_payload).await, + TierMutationRpcPhase::Commit => handle_commit(api, mutation_id, canonical_payload).await, + TierMutationRpcPhase::Abort => handle_abort(api, mutation_id, canonical_payload).await, + _ => Err(TierMutationPeerError::InvalidPayload( + "tier mutation rpc phase is unsupported".to_string(), + )), + } +} + +async fn handle_prepare( + api: Arc, + mutation_id: Uuid, + canonical_payload: &[u8], +) -> TierMutationPeerResult { + let intent = TierMutationIntent::decode(mutation_id, canonical_payload) + .map_err(|err| TierMutationPeerError::InvalidPayload(err.to_string()))?; + if intent.state != TierMutationIntentState::Prepared { + return Err(TierMutationPeerError::InvalidPayload( + "prepare intent must be in prepared state".to_string(), + )); + } + let tier_config_mgr = api.tier_config_mgr(); + + match save_tier_mutation_intent_record_if_absent(api.clone(), &intent).await { + Ok(()) => { + TierConfigMgr::apply_prepared_mutation_intent_block(&tier_config_mgr, &intent) + .await + .map_err(TierMutationPeerError::Runtime)?; + Ok(TierMutationPeerOutcome { + state: TierMutationPeerState::Prepared, + applied: true, + }) + } + Err(Error::PreconditionFailed) => { + let existing = load_tier_mutation_intent_record(api, mutation_id).await?; + if !same_mutation_identity(&existing, &intent) { + return Err(TierMutationPeerError::ConflictingIntent); + } + match existing.state { + TierMutationIntentState::Prepared => { + TierConfigMgr::apply_prepared_mutation_intent_block(&tier_config_mgr, &existing) + .await + .map_err(TierMutationPeerError::Runtime)?; + } + TierMutationIntentState::Committed | TierMutationIntentState::Aborted => { + TierConfigMgr::clear_prepared_mutation_intent_block(&tier_config_mgr, mutation_id).await; + } + } + Ok(TierMutationPeerOutcome { + state: peer_state_from_intent(existing.state), + applied: false, + }) + } + Err(err) => Err(err.into()), + } +} + +async fn handle_commit( + api: Arc, + mutation_id: Uuid, + canonical_payload: &[u8], +) -> TierMutationPeerResult { + let committed_config_etag = parse_commit_etag(canonical_payload)?; + let tier_config_mgr = api.tier_config_mgr(); + let (intent, applied) = advance_tier_mutation_intent_record_idempotent( + api, + mutation_id, + TierMutationIntentState::Committed, + Some(committed_config_etag), + ) + .await?; + if intent.state == TierMutationIntentState::Committed { + TierConfigMgr::clear_prepared_mutation_intent_block(&tier_config_mgr, mutation_id).await; + } + Ok(TierMutationPeerOutcome { + state: peer_state_from_intent(intent.state), + applied, + }) +} + +async fn handle_abort( + api: Arc, + mutation_id: Uuid, + canonical_payload: &[u8], +) -> TierMutationPeerResult { + if !canonical_payload.is_empty() { + return Err(TierMutationPeerError::InvalidPayload("abort payload must be empty".to_string())); + } + let tier_config_mgr = api.tier_config_mgr(); + let (intent, applied) = + advance_tier_mutation_intent_record_idempotent(api, mutation_id, TierMutationIntentState::Aborted, None).await?; + if intent.state == TierMutationIntentState::Aborted { + TierConfigMgr::clear_prepared_mutation_intent_block(&tier_config_mgr, mutation_id).await; + } + Ok(TierMutationPeerOutcome { + state: peer_state_from_intent(intent.state), + applied, + }) +} + +fn validate_peer_request_envelope( + protocol_version: u32, + mutation_id: Uuid, + canonical_payload: &[u8], +) -> TierMutationPeerResult<()> { + if protocol_version != TIER_MUTATION_RPC_PROTOCOL_VERSION { + return Err(TierMutationPeerError::UnsupportedProtocolVersion(protocol_version)); + } + if mutation_id.is_nil() { + return Err(TierMutationPeerError::NilMutationId); + } + if canonical_payload.len() > MAX_TIER_MUTATION_INTENT_SIZE { + return Err(TierMutationPeerError::PayloadTooLarge { + len: canonical_payload.len(), + max: MAX_TIER_MUTATION_INTENT_SIZE, + }); + } + Ok(()) +} + +fn parse_commit_etag(canonical_payload: &[u8]) -> TierMutationPeerResult { + if canonical_payload.len() > MAX_TIER_MUTATION_PEER_COMMIT_ETAG_SIZE { + return Err(TierMutationPeerError::PayloadTooLarge { + len: canonical_payload.len(), + max: MAX_TIER_MUTATION_PEER_COMMIT_ETAG_SIZE, + }); + } + let etag = std::str::from_utf8(canonical_payload) + .map_err(|err| TierMutationPeerError::InvalidPayload(err.to_string()))? + .trim(); + if etag.is_empty() { + return Err(TierMutationPeerError::InvalidPayload( + "commit payload must carry a committed config etag".to_string(), + )); + } + Ok(etag.to_string()) +} + +fn peer_state_from_intent(state: TierMutationIntentState) -> TierMutationPeerState { + match state { + TierMutationIntentState::Prepared => TierMutationPeerState::Prepared, + TierMutationIntentState::Committed => TierMutationPeerState::Committed, + TierMutationIntentState::Aborted => TierMutationPeerState::Aborted, + } +} + +fn same_mutation_identity(existing: &TierMutationIntent, expected: &TierMutationIntent) -> bool { + existing.mutation_id == expected.mutation_id + && existing.kind == expected.kind + && existing.old_config_etag == expected.old_config_etag + && existing.candidate_digest == expected.candidate_digest + && existing.affected_targets == expected.affected_targets + && existing.expires_at_unix_nanos == expected.expires_at_unix_nanos +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn peer_request_envelope_fails_closed_on_old_version_nil_id_and_large_payload() { + let mutation_id = Uuid::new_v4(); + assert!(matches!( + validate_peer_request_envelope(TIER_MUTATION_RPC_PROTOCOL_VERSION + 1, mutation_id, b"payload"), + Err(TierMutationPeerError::UnsupportedProtocolVersion(_)) + )); + assert!(matches!( + validate_peer_request_envelope(TIER_MUTATION_RPC_PROTOCOL_VERSION, Uuid::nil(), b"payload"), + Err(TierMutationPeerError::NilMutationId) + )); + + let oversized = vec![0; MAX_TIER_MUTATION_INTENT_SIZE + 1]; + assert!(matches!( + validate_peer_request_envelope(TIER_MUTATION_RPC_PROTOCOL_VERSION, mutation_id, &oversized), + Err(TierMutationPeerError::PayloadTooLarge { .. }) + )); + } + + #[test] + fn commit_payload_requires_small_non_empty_utf8_etag() { + assert_eq!( + parse_commit_etag(b" committed-etag ").expect("etag payload should parse"), + "committed-etag" + ); + assert!(matches!( + parse_commit_etag(b" "), + Err(TierMutationPeerError::InvalidPayload(message)) if message.contains("committed config etag") + )); + assert!(matches!( + parse_commit_etag(&[0xff]), + Err(TierMutationPeerError::InvalidPayload(message)) if message.contains("utf-8") + )); + + let oversized = vec![b'a'; MAX_TIER_MUTATION_PEER_COMMIT_ETAG_SIZE + 1]; + assert!(matches!( + parse_commit_etag(&oversized), + Err(TierMutationPeerError::PayloadTooLarge { .. }) + )); + } +} diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index d7da6dad9..fecd97a0d 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -553,6 +553,7 @@ mod tests { list_tier_mutation_intent_records, load_tier_mutation_intent_record, load_tier_mutation_intent_record_with_etag, save_tier_mutation_intent_record, save_tier_mutation_intent_record_if_current, }, + tier_mutation_peer::{TierMutationPeerError, TierMutationPeerState, handle_tier_mutation_peer_request}, warm_backend::WarmBackend, }, storage_api_contracts::{ @@ -572,6 +573,8 @@ mod tests { }; use http::HeaderMap; use rustfs_config::server_config::KVS; + #[cfg(feature = "test-util")] + use rustfs_protos::{TIER_MUTATION_RPC_PROTOCOL_VERSION, TierMutationRpcPhase}; use std::{ future::Future, io::Cursor, @@ -1459,6 +1462,225 @@ mod tests { assert_eq!(aborted_retry, aborted); } + #[cfg(feature = "test-util")] + fn tier_mutation_peer_test_intent( + mutation_id: uuid::Uuid, + tier_name: &str, + candidate_digest: [u8; 32], + ) -> TierMutationIntent { + 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, + 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, + } + } + + #[cfg(feature = "test-util")] + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn tier_mutation_peer_handler_applies_prepare_commit_and_abort_idempotently() { + 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-peer-handler", &[4])).await; + let mutation_id = uuid::Uuid::new_v4(); + let intent = tier_mutation_peer_test_intent(mutation_id, "COLD-A", [3; 32]); + let prepare_payload = intent.encode().expect("prepare intent should encode"); + register_mock_tier(&store.tier_config_mgr(), "COLD-A").await; + + let prepared = handle_tier_mutation_peer_request( + store.clone(), + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Prepare, + mutation_id, + &prepare_payload, + ) + .await + .expect("first prepare should create the peer intent"); + assert!(prepared.applied); + assert_eq!(prepared.state, TierMutationPeerState::Prepared); + let blocked = match TierConfigMgr::acquire_operation_lease(&store.tier_config_mgr(), "COLD-A").await { + Ok(_) => panic!("prepared peer mutation should block new tier operation leases"), + Err(err) => err, + }; + assert!( + blocked.message.contains("being replaced"), + "prepared peer mutation should reuse the existing blocked-tier error: {blocked}" + ); + + let retried_prepare = handle_tier_mutation_peer_request( + store.clone(), + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Prepare, + mutation_id, + &prepare_payload, + ) + .await + .expect("same prepare retry should be idempotent"); + assert!(!retried_prepare.applied); + assert_eq!(retried_prepare.state, TierMutationPeerState::Prepared); + let retried_blocked = match TierConfigMgr::acquire_operation_lease(&store.tier_config_mgr(), "COLD-A").await { + Ok(_) => panic!("prepared retry should keep blocking new tier operation leases"), + Err(err) => err, + }; + assert!( + retried_blocked.message.contains("being replaced"), + "prepared retry should keep the existing blocked-tier error: {retried_blocked}" + ); + + let committed = handle_tier_mutation_peer_request( + store.clone(), + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Commit, + mutation_id, + b"new-etag", + ) + .await + .expect("commit should advance the prepared peer intent"); + assert!(committed.applied); + assert_eq!(committed.state, TierMutationPeerState::Committed); + drop( + TierConfigMgr::acquire_operation_lease(&store.tier_config_mgr(), "COLD-A") + .await + .expect("committed peer mutation should clear the prepared runtime block"), + ); + + let retried_commit = handle_tier_mutation_peer_request( + store.clone(), + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Commit, + mutation_id, + b"new-etag", + ) + .await + .expect("same commit retry should be idempotent"); + assert!(!retried_commit.applied); + assert_eq!(retried_commit.state, TierMutationPeerState::Committed); + + let delayed_prepare_retry = handle_tier_mutation_peer_request( + store.clone(), + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Prepare, + mutation_id, + &prepare_payload, + ) + .await + .expect("delayed duplicate prepare should report the durable committed state"); + assert!(!delayed_prepare_retry.applied); + assert_eq!(delayed_prepare_retry.state, TierMutationPeerState::Committed); + drop( + TierConfigMgr::acquire_operation_lease(&store.tier_config_mgr(), "COLD-A") + .await + .expect("delayed committed prepare retry must not recreate a runtime block"), + ); + + let loaded = load_tier_mutation_intent_record(store.clone(), mutation_id) + .await + .expect("committed peer intent should remain durable"); + assert_eq!(loaded.state, TierMutationIntentState::Committed); + assert_eq!(loaded.committed_config_etag.as_deref(), Some("new-etag")); + + let abort_id = uuid::Uuid::new_v4(); + let abort_intent = tier_mutation_peer_test_intent(abort_id, "COLD-B", [4; 32]); + let abort_prepare_payload = abort_intent.encode().expect("abort prepare intent should encode"); + register_mock_tier(&store.tier_config_mgr(), "COLD-B").await; + handle_tier_mutation_peer_request( + store.clone(), + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Prepare, + abort_id, + &abort_prepare_payload, + ) + .await + .expect("abort target prepare should create the peer intent"); + let abort_blocked = match TierConfigMgr::acquire_operation_lease(&store.tier_config_mgr(), "COLD-B").await { + Ok(_) => panic!("abort target prepare should block new tier operation leases"), + Err(err) => err, + }; + assert!( + abort_blocked.message.contains("being replaced"), + "abort target prepare should reuse the existing blocked-tier error: {abort_blocked}" + ); + + let aborted = handle_tier_mutation_peer_request( + store.clone(), + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Abort, + abort_id, + b"", + ) + .await + .expect("abort should advance the prepared peer intent"); + assert!(aborted.applied); + assert_eq!(aborted.state, TierMutationPeerState::Aborted); + drop( + TierConfigMgr::acquire_operation_lease(&store.tier_config_mgr(), "COLD-B") + .await + .expect("aborted peer mutation should clear the prepared runtime block"), + ); + + let retried_abort = handle_tier_mutation_peer_request( + store, + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Abort, + abort_id, + b"", + ) + .await + .expect("same abort retry should be idempotent"); + assert!(!retried_abort.applied); + assert_eq!(retried_abort.state, TierMutationPeerState::Aborted); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn tier_mutation_peer_handler_rejects_conflicting_prepare_without_overwrite() { + 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-peer-conflict", &[4])).await; + let mutation_id = uuid::Uuid::new_v4(); + let intent = tier_mutation_peer_test_intent(mutation_id, "COLD-A", [3; 32]); + let prepare_payload = intent.encode().expect("prepare intent should encode"); + handle_tier_mutation_peer_request( + store.clone(), + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Prepare, + mutation_id, + &prepare_payload, + ) + .await + .expect("first prepare should create the peer intent"); + + let conflicting = tier_mutation_peer_test_intent(mutation_id, "COLD-A", [4; 32]); + let conflicting_payload = conflicting.encode().expect("conflicting intent should encode"); + let conflict = handle_tier_mutation_peer_request( + store.clone(), + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Prepare, + mutation_id, + &conflicting_payload, + ) + .await + .expect_err("conflicting prepare must fail closed"); + assert!(matches!(conflict, TierMutationPeerError::ConflictingIntent)); + + let loaded = load_tier_mutation_intent_record(store, mutation_id) + .await + .expect("conflicting prepare must not overwrite the first record"); + assert_eq!(loaded.candidate_digest, [3; 32]); + assert_eq!(loaded.state, TierMutationIntentState::Prepared); + } + #[cfg(feature = "test-util")] #[tokio::test] #[serial_test::serial(storage_class_env)] diff --git a/crates/protos/src/generated/proto_gen/node_service.rs b/crates/protos/src/generated/proto_gen/node_service.rs index 0a0f1eb77..d22104b1d 100644 --- a/crates/protos/src/generated/proto_gen/node_service.rs +++ b/crates/protos/src/generated/proto_gen/node_service.rs @@ -1223,6 +1223,46 @@ pub struct LoadTransitionTierConfigResponse { #[prost(string, optional, tag = "2")] pub error_info: ::core::option::Option<::prost::alloc::string::String>, } +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct TierMutationPrepareRequest { + #[prost(uint32, tag = "1")] + pub version: u32, + #[prost(string, tag = "2")] + pub mutation_id: ::prost::alloc::string::String, + #[prost(bytes = "bytes", tag = "3")] + pub canonical_payload: ::prost::bytes::Bytes, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct TierMutationCommitRequest { + #[prost(uint32, tag = "1")] + pub version: u32, + #[prost(string, tag = "2")] + pub mutation_id: ::prost::alloc::string::String, + #[prost(bytes = "bytes", tag = "3")] + pub canonical_payload: ::prost::bytes::Bytes, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct TierMutationAbortRequest { + #[prost(uint32, tag = "1")] + pub version: u32, + #[prost(string, tag = "2")] + pub mutation_id: ::prost::alloc::string::String, + #[prost(bytes = "bytes", tag = "3")] + pub canonical_payload: ::prost::bytes::Bytes, +} +#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] +pub struct TierMutationControlResponse { + #[prost(bool, tag = "1")] + pub success: bool, + #[prost(enumeration = "TierMutationPeerState", tag = "2")] + pub state: i32, + #[prost(bool, tag = "3")] + pub applied: bool, + #[prost(string, optional, tag = "4")] + pub error_info: ::core::option::Option<::prost::alloc::string::String>, + #[prost(bytes = "bytes", tag = "5")] + pub response_proof: ::prost::bytes::Bytes, +} #[derive(Clone, Copy, PartialEq, Eq, Hash, ::prost::Message)] pub struct GetLiveEventsRequest { #[prost(uint64, tag = "1")] @@ -1243,6 +1283,38 @@ pub struct GetLiveEventsResponse { #[prost(string, optional, tag = "5")] pub error_info: ::core::option::Option<::prost::alloc::string::String>, } +#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)] +#[repr(i32)] +pub enum TierMutationPeerState { + Unspecified = 0, + Prepared = 1, + Committed = 2, + Aborted = 3, +} +impl TierMutationPeerState { + /// String value of the enum field names used in the ProtoBuf definition. + /// + /// The values are not transformed in any way and thus are considered stable + /// (if the ProtoBuf definition does not change) and safe for programmatic use. + pub fn as_str_name(&self) -> &'static str { + match self { + Self::Unspecified => "TIER_MUTATION_PEER_STATE_UNSPECIFIED", + Self::Prepared => "TIER_MUTATION_PEER_STATE_PREPARED", + Self::Committed => "TIER_MUTATION_PEER_STATE_COMMITTED", + Self::Aborted => "TIER_MUTATION_PEER_STATE_ABORTED", + } + } + /// Creates an enum from field names used in the ProtoBuf definition. + pub fn from_str_name(value: &str) -> ::core::option::Option { + match value { + "TIER_MUTATION_PEER_STATE_UNSPECIFIED" => Some(Self::Unspecified), + "TIER_MUTATION_PEER_STATE_PREPARED" => Some(Self::Prepared), + "TIER_MUTATION_PEER_STATE_COMMITTED" => Some(Self::Committed), + "TIER_MUTATION_PEER_STATE_ABORTED" => Some(Self::Aborted), + _ => None, + } + } +} /// Generated client implementations. pub mod node_service_client { #![allow(unused_variables, dead_code, missing_docs, clippy::wildcard_imports, clippy::let_unit_value)] @@ -5590,3 +5662,334 @@ pub mod heal_control_service_server { const NAME: &'static str = SERVICE_NAME; } } +/// Generated client implementations. +pub mod tier_mutation_control_service_client { + #![allow(unused_variables, dead_code, missing_docs, clippy::wildcard_imports, clippy::let_unit_value)] + use tonic::codegen::http::Uri; + use tonic::codegen::*; + #[derive(Debug, Clone)] + pub struct TierMutationControlServiceClient { + inner: tonic::client::Grpc, + } + impl TierMutationControlServiceClient { + /// Attempt to create a new client by connecting to a given endpoint. + pub async fn connect(dst: D) -> Result + where + D: TryInto, + D::Error: Into, + { + let conn = tonic::transport::Endpoint::new(dst)?.connect().await?; + Ok(Self::new(conn)) + } + } + impl TierMutationControlServiceClient + where + T: tonic::client::GrpcService, + T::Error: Into, + T::ResponseBody: Body + std::marker::Send + 'static, + ::Error: Into + std::marker::Send, + { + pub fn new(inner: T) -> Self { + let inner = tonic::client::Grpc::new(inner); + Self { inner } + } + pub fn with_origin(inner: T, origin: Uri) -> Self { + let inner = tonic::client::Grpc::with_origin(inner, origin); + Self { inner } + } + pub fn with_interceptor(inner: T, interceptor: F) -> TierMutationControlServiceClient> + where + F: tonic::service::Interceptor, + T::ResponseBody: Default, + T: tonic::codegen::Service< + http::Request, + Response = http::Response<>::ResponseBody>, + >, + >>::Error: + Into + std::marker::Send + std::marker::Sync, + { + TierMutationControlServiceClient::new(InterceptedService::new(inner, interceptor)) + } + /// Compress requests with the given encoding. + /// + /// This requires the server to support it otherwise it might respond with an + /// error. + #[must_use] + pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self { + self.inner = self.inner.send_compressed(encoding); + self + } + /// Enable decompressing responses. + #[must_use] + pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self { + self.inner = self.inner.accept_compressed(encoding); + self + } + /// Limits the maximum size of a decoded message. + /// + /// Default: `4MB` + #[must_use] + pub fn max_decoding_message_size(mut self, limit: usize) -> Self { + self.inner = self.inner.max_decoding_message_size(limit); + self + } + /// Limits the maximum size of an encoded message. + /// + /// Default: `usize::MAX` + #[must_use] + pub fn max_encoding_message_size(mut self, limit: usize) -> Self { + self.inner = self.inner.max_encoding_message_size(limit); + self + } + pub async fn prepare_tier_mutation( + &mut self, + request: impl tonic::IntoRequest, + ) -> std::result::Result, tonic::Status> { + self.inner + .ready() + .await + .map_err(|e| tonic::Status::unknown(format!("Service was not ready: {}", e.into())))?; + let codec = tonic_prost::ProstCodec::default(); + let path = http::uri::PathAndQuery::from_static("/node_service.TierMutationControlService/PrepareTierMutation"); + let mut req = request.into_request(); + req.extensions_mut() + .insert(GrpcMethod::new("node_service.TierMutationControlService", "PrepareTierMutation")); + self.inner.unary(req, path, codec).await + } + pub async fn commit_tier_mutation( + &mut self, + request: impl tonic::IntoRequest, + ) -> std::result::Result, tonic::Status> { + self.inner + .ready() + .await + .map_err(|e| tonic::Status::unknown(format!("Service was not ready: {}", e.into())))?; + let codec = tonic_prost::ProstCodec::default(); + let path = http::uri::PathAndQuery::from_static("/node_service.TierMutationControlService/CommitTierMutation"); + let mut req = request.into_request(); + req.extensions_mut() + .insert(GrpcMethod::new("node_service.TierMutationControlService", "CommitTierMutation")); + self.inner.unary(req, path, codec).await + } + pub async fn abort_tier_mutation( + &mut self, + request: impl tonic::IntoRequest, + ) -> std::result::Result, tonic::Status> { + self.inner + .ready() + .await + .map_err(|e| tonic::Status::unknown(format!("Service was not ready: {}", e.into())))?; + let codec = tonic_prost::ProstCodec::default(); + let path = http::uri::PathAndQuery::from_static("/node_service.TierMutationControlService/AbortTierMutation"); + let mut req = request.into_request(); + req.extensions_mut() + .insert(GrpcMethod::new("node_service.TierMutationControlService", "AbortTierMutation")); + self.inner.unary(req, path, codec).await + } + } +} +/// Generated server implementations. +pub mod tier_mutation_control_service_server { + #![allow(unused_variables, dead_code, missing_docs, clippy::wildcard_imports, clippy::let_unit_value)] + use tonic::codegen::*; + /// Generated trait containing gRPC methods that should be implemented for use with TierMutationControlServiceServer. + #[async_trait] + pub trait TierMutationControlService: std::marker::Send + std::marker::Sync + 'static { + async fn prepare_tier_mutation( + &self, + request: tonic::Request, + ) -> std::result::Result, tonic::Status>; + async fn commit_tier_mutation( + &self, + request: tonic::Request, + ) -> std::result::Result, tonic::Status>; + async fn abort_tier_mutation( + &self, + request: tonic::Request, + ) -> std::result::Result, tonic::Status>; + } + #[derive(Debug)] + pub struct TierMutationControlServiceServer { + inner: Arc, + accept_compression_encodings: EnabledCompressionEncodings, + send_compression_encodings: EnabledCompressionEncodings, + max_decoding_message_size: Option, + max_encoding_message_size: Option, + } + impl TierMutationControlServiceServer { + pub fn new(inner: T) -> Self { + Self::from_arc(Arc::new(inner)) + } + pub fn from_arc(inner: Arc) -> Self { + Self { + inner, + accept_compression_encodings: Default::default(), + send_compression_encodings: Default::default(), + max_decoding_message_size: None, + max_encoding_message_size: None, + } + } + pub fn with_interceptor(inner: T, interceptor: F) -> InterceptedService + where + F: tonic::service::Interceptor, + { + InterceptedService::new(Self::new(inner), interceptor) + } + /// Enable decompressing requests with the given encoding. + #[must_use] + pub fn accept_compressed(mut self, encoding: CompressionEncoding) -> Self { + self.accept_compression_encodings.enable(encoding); + self + } + /// Compress responses with the given encoding, if the client supports it. + #[must_use] + pub fn send_compressed(mut self, encoding: CompressionEncoding) -> Self { + self.send_compression_encodings.enable(encoding); + self + } + /// Limits the maximum size of a decoded message. + /// + /// Default: `4MB` + #[must_use] + pub fn max_decoding_message_size(mut self, limit: usize) -> Self { + self.max_decoding_message_size = Some(limit); + self + } + /// Limits the maximum size of an encoded message. + /// + /// Default: `usize::MAX` + #[must_use] + pub fn max_encoding_message_size(mut self, limit: usize) -> Self { + self.max_encoding_message_size = Some(limit); + self + } + } + impl tonic::codegen::Service> for TierMutationControlServiceServer + where + T: TierMutationControlService, + B: Body + std::marker::Send + 'static, + B::Error: Into + std::marker::Send + 'static, + { + type Response = http::Response; + type Error = std::convert::Infallible; + type Future = BoxFuture; + fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) + } + fn call(&mut self, req: http::Request) -> Self::Future { + match req.uri().path() { + "/node_service.TierMutationControlService/PrepareTierMutation" => { + #[allow(non_camel_case_types)] + struct PrepareTierMutationSvc(pub Arc); + impl tonic::server::UnaryService for PrepareTierMutationSvc { + type Response = super::TierMutationControlResponse; + type Future = BoxFuture, tonic::Status>; + fn call(&mut self, request: tonic::Request) -> Self::Future { + let inner = Arc::clone(&self.0); + let fut = + async move { ::prepare_tier_mutation(&inner, request).await }; + Box::pin(fut) + } + } + let accept_compression_encodings = self.accept_compression_encodings; + let send_compression_encodings = self.send_compression_encodings; + let max_decoding_message_size = self.max_decoding_message_size; + let max_encoding_message_size = self.max_encoding_message_size; + let inner = self.inner.clone(); + let fut = async move { + let method = PrepareTierMutationSvc(inner); + let codec = tonic_prost::ProstCodec::default(); + let mut grpc = tonic::server::Grpc::new(codec) + .apply_compression_config(accept_compression_encodings, send_compression_encodings) + .apply_max_message_size_config(max_decoding_message_size, max_encoding_message_size); + let res = grpc.unary(method, req).await; + Ok(res) + }; + Box::pin(fut) + } + "/node_service.TierMutationControlService/CommitTierMutation" => { + #[allow(non_camel_case_types)] + struct CommitTierMutationSvc(pub Arc); + impl tonic::server::UnaryService for CommitTierMutationSvc { + type Response = super::TierMutationControlResponse; + type Future = BoxFuture, tonic::Status>; + fn call(&mut self, request: tonic::Request) -> Self::Future { + let inner = Arc::clone(&self.0); + let fut = + async move { ::commit_tier_mutation(&inner, request).await }; + Box::pin(fut) + } + } + let accept_compression_encodings = self.accept_compression_encodings; + let send_compression_encodings = self.send_compression_encodings; + let max_decoding_message_size = self.max_decoding_message_size; + let max_encoding_message_size = self.max_encoding_message_size; + let inner = self.inner.clone(); + let fut = async move { + let method = CommitTierMutationSvc(inner); + let codec = tonic_prost::ProstCodec::default(); + let mut grpc = tonic::server::Grpc::new(codec) + .apply_compression_config(accept_compression_encodings, send_compression_encodings) + .apply_max_message_size_config(max_decoding_message_size, max_encoding_message_size); + let res = grpc.unary(method, req).await; + Ok(res) + }; + Box::pin(fut) + } + "/node_service.TierMutationControlService/AbortTierMutation" => { + #[allow(non_camel_case_types)] + struct AbortTierMutationSvc(pub Arc); + impl tonic::server::UnaryService for AbortTierMutationSvc { + type Response = super::TierMutationControlResponse; + type Future = BoxFuture, tonic::Status>; + fn call(&mut self, request: tonic::Request) -> Self::Future { + let inner = Arc::clone(&self.0); + let fut = + async move { ::abort_tier_mutation(&inner, request).await }; + Box::pin(fut) + } + } + let accept_compression_encodings = self.accept_compression_encodings; + let send_compression_encodings = self.send_compression_encodings; + let max_decoding_message_size = self.max_decoding_message_size; + let max_encoding_message_size = self.max_encoding_message_size; + let inner = self.inner.clone(); + let fut = async move { + let method = AbortTierMutationSvc(inner); + let codec = tonic_prost::ProstCodec::default(); + let mut grpc = tonic::server::Grpc::new(codec) + .apply_compression_config(accept_compression_encodings, send_compression_encodings) + .apply_max_message_size_config(max_decoding_message_size, max_encoding_message_size); + let res = grpc.unary(method, req).await; + Ok(res) + }; + Box::pin(fut) + } + _ => Box::pin(async move { + let mut response = http::Response::new(tonic::body::Body::default()); + let headers = response.headers_mut(); + headers.insert(tonic::Status::GRPC_STATUS, (tonic::Code::Unimplemented as i32).into()); + headers.insert(http::header::CONTENT_TYPE, tonic::metadata::GRPC_CONTENT_TYPE); + Ok(response) + }), + } + } + } + impl Clone for TierMutationControlServiceServer { + fn clone(&self) -> Self { + let inner = self.inner.clone(); + Self { + inner, + accept_compression_encodings: self.accept_compression_encodings, + send_compression_encodings: self.send_compression_encodings, + max_decoding_message_size: self.max_decoding_message_size, + max_encoding_message_size: self.max_encoding_message_size, + } + } + } + /// Generated gRPC service name + pub const SERVICE_NAME: &str = "node_service.TierMutationControlService"; + impl tonic::server::NamedService for TierMutationControlServiceServer { + const NAME: &'static str = SERVICE_NAME; + } +} diff --git a/crates/protos/src/lib.rs b/crates/protos/src/lib.rs index 2cf2db8d3..fc34ae66e 100644 --- a/crates/protos/src/lib.rs +++ b/crates/protos/src/lib.rs @@ -35,6 +35,7 @@ use tonic::{ transport::{Certificate, Channel, ClientTlsConfig, Endpoint}, }; use tracing::{debug, info, warn}; +use uuid::Uuid; // Type alias for the complex client type pub type NodeServiceClientType = NodeServiceClient< @@ -165,6 +166,9 @@ pub fn internode_rpc_max_message_size() -> usize { pub const HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE: usize = heal_control::RESULT_MAX_SIZE + 1024; pub const HEAL_CONTROL_PROTOCOL_VERSION: u32 = 2; pub const HEAL_CONTROL_CAPABILITY_PROBE_PREFIX: &[u8] = b"rustfs-heal-control-capability-v2\0"; +pub const TIER_MUTATION_RPC_MAX_PREPARE_PAYLOAD_SIZE: usize = 64 * 1024; +pub const TIER_MUTATION_RPC_MAX_COMMIT_PAYLOAD_SIZE: usize = 1024; +pub const TIER_MUTATION_RPC_MAX_MESSAGE_SIZE: usize = TIER_MUTATION_RPC_MAX_PREPARE_PAYLOAD_SIZE + 4096; pub fn heal_control_coordinator_epoch(topology_fingerprint: &str) -> Result { let prefix = topology_fingerprint @@ -251,6 +255,100 @@ pub fn canonical_heal_control_response_body( Ok(body) } +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +#[non_exhaustive] +pub enum TierMutationRpcPhase { + Prepare, + Commit, + Abort, +} + +impl TierMutationRpcPhase { + fn as_wire_str(self) -> &'static str { + match self { + Self::Prepare => "prepare", + Self::Commit => "commit", + Self::Abort => "abort", + } + } +} + +pub const TIER_MUTATION_RPC_PROTOCOL_VERSION: u32 = 1; + +pub fn canonical_tier_mutation_rpc_body( + version: u32, + phase: TierMutationRpcPhase, + mutation_id: Uuid, + canonical_payload: &[u8], +) -> Result, std::num::TryFromIntError> { + const DOMAIN: &[u8] = b"rustfs-tier-mutation-rpc-v1\0"; + + let phase = phase.as_wire_str().as_bytes(); + let mutation_id = mutation_id.as_bytes(); + let mut body = Vec::with_capacity(DOMAIN.len() + 4 + 8 + phase.len() + mutation_id.len() + 8 + canonical_payload.len()); + body.extend_from_slice(DOMAIN); + body.extend_from_slice(&version.to_be_bytes()); + body.extend_from_slice(&u64::try_from(phase.len())?.to_be_bytes()); + body.extend_from_slice(phase); + body.extend_from_slice(mutation_id); + body.extend_from_slice(&u64::try_from(canonical_payload.len())?.to_be_bytes()); + body.extend_from_slice(canonical_payload); + Ok(body) +} + +pub struct TierMutationRpcResponseProofInput<'a> { + pub version: u32, + pub phase: TierMutationRpcPhase, + pub mutation_id: Uuid, + pub canonical_payload: &'a [u8], + pub success: bool, + pub state: i32, + pub applied: bool, + pub error_info: Option<&'a str>, +} + +pub fn canonical_tier_mutation_rpc_response_body( + input: TierMutationRpcResponseProofInput<'_>, +) -> Result, std::num::TryFromIntError> { + const DOMAIN: &[u8] = b"rustfs-tier-mutation-rpc-response-v1\0"; + + let phase = input.phase.as_wire_str().as_bytes(); + let mutation_id = input.mutation_id.as_bytes(); + let error_info = input.error_info.map(str::as_bytes); + let error_info_len = error_info.map_or(0, <[u8]>::len); + let mut body = Vec::with_capacity( + DOMAIN.len() + + 4 + + 8 + + phase.len() + + mutation_id.len() + + 8 + + input.canonical_payload.len() + + 1 + + 4 + + 1 + + 1 + + 8 + + error_info_len, + ); + body.extend_from_slice(DOMAIN); + body.extend_from_slice(&input.version.to_be_bytes()); + body.extend_from_slice(&u64::try_from(phase.len())?.to_be_bytes()); + body.extend_from_slice(phase); + body.extend_from_slice(mutation_id); + body.extend_from_slice(&u64::try_from(input.canonical_payload.len())?.to_be_bytes()); + body.extend_from_slice(input.canonical_payload); + body.push(u8::from(input.success)); + body.extend_from_slice(&input.state.to_be_bytes()); + body.push(u8::from(input.applied)); + body.push(u8::from(error_info.is_some())); + body.extend_from_slice(&u64::try_from(error_info_len)?.to_be_bytes()); + if let Some(error_info) = error_info { + body.extend_from_slice(error_info); + } + Ok(body) +} + #[cfg(test)] mod heal_control_tests { use super::{ @@ -363,6 +461,179 @@ mod heal_control_tests { } } +#[cfg(test)] +mod tier_mutation_rpc_tests { + use super::{ + TIER_MUTATION_RPC_PROTOCOL_VERSION, TierMutationRpcPhase, TierMutationRpcResponseProofInput, + canonical_tier_mutation_rpc_body, canonical_tier_mutation_rpc_response_body, + }; + use crate::proto_gen::node_service::TierMutationPeerState; + use uuid::uuid; + + #[test] + fn canonical_tier_mutation_body_binds_phase_id_and_payload() { + let mutation_id = uuid!("12345678-1234-5678-9abc-def012345678"); + let payload = b"canonical-intent-record"; + let baseline = canonical_tier_mutation_rpc_body( + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Prepare, + mutation_id, + payload, + ) + .expect("small mutation body should encode"); + let mut golden = b"rustfs-tier-mutation-rpc-v1\0".to_vec(); + golden.extend_from_slice(&1_u32.to_be_bytes()); + golden.extend_from_slice(&7_u64.to_be_bytes()); + golden.extend_from_slice(b"prepare"); + golden.extend_from_slice(mutation_id.as_bytes()); + golden.extend_from_slice(&u64::try_from(payload.len()).expect("payload length should fit").to_be_bytes()); + golden.extend_from_slice(payload); + assert_eq!(baseline, golden); + + assert_ne!( + baseline, + canonical_tier_mutation_rpc_body(2, TierMutationRpcPhase::Prepare, mutation_id, payload) + .expect("small mutation body should encode") + ); + assert_ne!( + baseline, + canonical_tier_mutation_rpc_body( + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Commit, + mutation_id, + payload, + ) + .expect("small mutation body should encode") + ); + assert_ne!( + baseline, + canonical_tier_mutation_rpc_body( + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Prepare, + uuid!("22345678-1234-5678-9abc-def012345678"), + payload, + ) + .expect("small mutation body should encode") + ); + assert_ne!( + baseline, + canonical_tier_mutation_rpc_body( + TIER_MUTATION_RPC_PROTOCOL_VERSION, + TierMutationRpcPhase::Prepare, + mutation_id, + b"canonical-intent-record-tampered", + ) + .expect("small mutation body should encode") + ); + } + + #[test] + fn canonical_tier_mutation_response_binds_request_state_and_error() { + let mutation_id = uuid!("12345678-1234-5678-9abc-def012345678"); + let payload = b"canonical-intent-record"; + let baseline = canonical_tier_mutation_rpc_response_body(TierMutationRpcResponseProofInput { + version: TIER_MUTATION_RPC_PROTOCOL_VERSION, + phase: TierMutationRpcPhase::Prepare, + mutation_id, + canonical_payload: payload, + success: true, + state: TierMutationPeerState::Prepared as i32, + applied: true, + error_info: None, + }) + .expect("small mutation response should encode"); + + let cases = [ + TierMutationRpcResponseProofInput { + version: 2, + phase: TierMutationRpcPhase::Prepare, + mutation_id, + canonical_payload: payload, + success: true, + state: TierMutationPeerState::Prepared as i32, + applied: true, + error_info: None, + }, + TierMutationRpcResponseProofInput { + version: TIER_MUTATION_RPC_PROTOCOL_VERSION, + phase: TierMutationRpcPhase::Commit, + mutation_id, + canonical_payload: payload, + success: true, + state: TierMutationPeerState::Prepared as i32, + applied: true, + error_info: None, + }, + TierMutationRpcResponseProofInput { + version: TIER_MUTATION_RPC_PROTOCOL_VERSION, + phase: TierMutationRpcPhase::Prepare, + mutation_id: uuid!("22345678-1234-5678-9abc-def012345678"), + canonical_payload: payload, + success: true, + state: TierMutationPeerState::Prepared as i32, + applied: true, + error_info: None, + }, + TierMutationRpcResponseProofInput { + version: TIER_MUTATION_RPC_PROTOCOL_VERSION, + phase: TierMutationRpcPhase::Prepare, + mutation_id, + canonical_payload: b"tampered-intent-record", + success: true, + state: TierMutationPeerState::Prepared as i32, + applied: true, + error_info: None, + }, + TierMutationRpcResponseProofInput { + version: TIER_MUTATION_RPC_PROTOCOL_VERSION, + phase: TierMutationRpcPhase::Prepare, + mutation_id, + canonical_payload: payload, + success: false, + state: TierMutationPeerState::Prepared as i32, + applied: true, + error_info: None, + }, + TierMutationRpcResponseProofInput { + version: TIER_MUTATION_RPC_PROTOCOL_VERSION, + phase: TierMutationRpcPhase::Prepare, + mutation_id, + canonical_payload: payload, + success: true, + state: TierMutationPeerState::Committed as i32, + applied: true, + error_info: None, + }, + TierMutationRpcResponseProofInput { + version: TIER_MUTATION_RPC_PROTOCOL_VERSION, + phase: TierMutationRpcPhase::Prepare, + mutation_id, + canonical_payload: payload, + success: true, + state: TierMutationPeerState::Prepared as i32, + applied: false, + error_info: None, + }, + TierMutationRpcResponseProofInput { + version: TIER_MUTATION_RPC_PROTOCOL_VERSION, + phase: TierMutationRpcPhase::Prepare, + mutation_id, + canonical_payload: payload, + success: true, + state: TierMutationPeerState::Prepared as i32, + applied: true, + error_info: Some("error"), + }, + ]; + for case in cases { + assert_ne!( + baseline, + canonical_tier_mutation_rpc_response_body(case).expect("small mutation response should encode") + ); + } + } +} + /// Whether internode metadata RPCs should send only the msgpack `_bin` payloads and leave the JSON /// compatibility strings empty (grpc-optimization P2-1). Shared by the client (`remote_disk`) and /// server (`node_service`) send paths. Defaults to `false` (dual-write); see diff --git a/crates/protos/src/node.proto b/crates/protos/src/node.proto index 14b5eb808..918f92c79 100644 --- a/crates/protos/src/node.proto +++ b/crates/protos/src/node.proto @@ -865,6 +865,39 @@ message LoadTransitionTierConfigResponse { optional string error_info = 2; } +message TierMutationPrepareRequest { + uint32 version = 1; + string mutation_id = 2; + bytes canonical_payload = 3; +} + +message TierMutationCommitRequest { + uint32 version = 1; + string mutation_id = 2; + bytes canonical_payload = 3; +} + +message TierMutationAbortRequest { + uint32 version = 1; + string mutation_id = 2; + bytes canonical_payload = 3; +} + +enum TierMutationPeerState { + TIER_MUTATION_PEER_STATE_UNSPECIFIED = 0; + TIER_MUTATION_PEER_STATE_PREPARED = 1; + TIER_MUTATION_PEER_STATE_COMMITTED = 2; + TIER_MUTATION_PEER_STATE_ABORTED = 3; +} + +message TierMutationControlResponse { + bool success = 1; + TierMutationPeerState state = 2; + bool applied = 3; + optional string error_info = 4; + bytes response_proof = 5; +} + message GetLiveEventsRequest { uint64 after_sequence = 1; uint32 limit = 2; @@ -983,3 +1016,9 @@ service NodeService { service HealControlService { rpc HealControl(HealControlRequest) returns (HealControlResponse) {}; } + +service TierMutationControlService { + rpc PrepareTierMutation(TierMutationPrepareRequest) returns (TierMutationControlResponse) {}; + rpc CommitTierMutation(TierMutationCommitRequest) returns (TierMutationControlResponse) {}; + rpc AbortTierMutation(TierMutationAbortRequest) returns (TierMutationControlResponse) {}; +} diff --git a/rustfs/src/server/http.rs b/rustfs/src/server/http.rs index 1b48e657f..d072c2c9e 100644 --- a/rustfs/src/server/http.rs +++ b/rustfs/src/server/http.rs @@ -57,6 +57,7 @@ use rustfs_keystone::KeystoneAuthLayer; use rustfs_protocols::SwiftService; use rustfs_protos::proto_gen::node_service::{ heal_control_service_server::HealControlServiceServer, node_service_server::NodeServiceServer, + tier_mutation_control_service_server::TierMutationControlServiceServer, }; use rustfs_trusted_proxies::ClientInfo; use rustfs_utils::net::parse_and_resolve_address; @@ -128,6 +129,9 @@ const EVENT_PEER_ADDR_UNAVAILABLE: &str = "peer_addr_unavailable"; const EVENT_RPC_SIGNATURE_VERIFICATION_FAILED: &str = "rpc_signature_verification_failed"; const EVENT_GRPC_TRACE_CONTEXT_PROPAGATION_FAILED: &str = "grpc_trace_context_propagation_failed"; const HEAL_CONTROL_TONIC_RPC_PATH: &str = "/node_service.HealControlService/HealControl"; +const TIER_MUTATION_PREPARE_TONIC_RPC_PATH: &str = "/node_service.TierMutationControlService/PrepareTierMutation"; +const TIER_MUTATION_COMMIT_TONIC_RPC_PATH: &str = "/node_service.TierMutationControlService/CommitTierMutation"; +const TIER_MUTATION_ABORT_TONIC_RPC_PATH: &str = "/node_service.TierMutationControlService/AbortTierMutation"; static ACTIVE_HTTP_REQUESTS: AtomicU64 = AtomicU64::new(0); @@ -1322,7 +1326,19 @@ fn process_connection( .max_encoding_message_size(heal_control_max_message_size), check_auth, ); - let rpc_service = RpcRequestPathService::new(Routes::new(node_service).add_service(heal_control_service).prepare()); + let tier_mutation_control_max_message_size = rustfs_protos::TIER_MUTATION_RPC_MAX_MESSAGE_SIZE; + let tier_mutation_control_service = InterceptedService::new( + TierMutationControlServiceServer::new(storage::tonic_service::make_tier_mutation_control_server()) + .max_decoding_message_size(tier_mutation_control_max_message_size) + .max_encoding_message_size(tier_mutation_control_max_message_size), + check_auth, + ); + let rpc_service = RpcRequestPathService::new( + Routes::new(node_service) + .add_service(heal_control_service) + .add_service(tier_mutation_control_service) + .prepare(), + ); #[cfg(feature = "swift")] let http_service = SwiftService::new(true, None, s3_service); @@ -1861,6 +1877,9 @@ fn check_auth(req: Request<()>) -> std::result::Result, Status> { .strip_prefix(TONIC_RPC_PREFIX) .and_then(|suffix| suffix.strip_prefix('/')) .or_else(|| (target.uri.path() == HEAL_CONTROL_TONIC_RPC_PATH).then_some("HealControl")) + .or_else(|| (target.uri.path() == TIER_MUTATION_PREPARE_TONIC_RPC_PATH).then_some("PrepareTierMutation")) + .or_else(|| (target.uri.path() == TIER_MUTATION_COMMIT_TONIC_RPC_PATH).then_some("CommitTierMutation")) + .or_else(|| (target.uri.path() == TIER_MUTATION_ABORT_TONIC_RPC_PATH).then_some("AbortTierMutation")) .filter(|method| !method.is_empty() && !method.contains('/')) .ok_or_else(|| Status::unauthenticated("Invalid RPC request path"))?; debug_assert!(!rpc_method.is_empty()); @@ -2279,6 +2298,23 @@ mod tests { }); assert!(check_auth(heal_request).is_ok(), "heal control service path should authenticate"); + let tier_headers = storage::gen_tonic_signature_headers( + "127.0.0.1:9000", + "node_service.TierMutationControlService", + "PrepareTierMutation", + None, + ) + .expect("tier mutation auth headers should build"); + let mut tier_request = Request::new(()); + tier_request.metadata_mut().as_mut().extend(tier_headers); + tier_request.extensions_mut().insert(RpcRequestTarget { + uri: TIER_MUTATION_PREPARE_TONIC_RPC_PATH + .parse() + .expect("tier mutation path should parse"), + method: Method::POST, + }); + assert!(check_auth(tier_request).is_ok(), "tier mutation control service path should authenticate"); + let replay_headers = storage::gen_tonic_signature_headers("127.0.0.1:9000", "node_service.NodeService", "Ping", None) .expect("node service auth headers should build"); let mut cross_service_replay = Request::new(()); @@ -2292,6 +2328,21 @@ mod tests { "node service signature must not replay to heal control" ); + let replay_headers = storage::gen_tonic_signature_headers("127.0.0.1:9000", "node_service.NodeService", "Ping", None) + .expect("node service auth headers should build"); + let mut cross_service_replay = Request::new(()); + cross_service_replay.metadata_mut().as_mut().extend(replay_headers); + cross_service_replay.extensions_mut().insert(RpcRequestTarget { + uri: TIER_MUTATION_PREPARE_TONIC_RPC_PATH + .parse() + .expect("tier mutation path should parse"), + method: Method::POST, + }); + assert!( + check_auth(cross_service_replay).is_err(), + "node service signature must not replay to tier mutation control" + ); + rustfs_common::set_global_local_node_name("127.0.0.1:9001").await; let mut replay_to_other_node = Request::new(()); replay_to_other_node.metadata_mut().as_mut().extend(headers.clone()); diff --git a/rustfs/src/storage/rpc/mod.rs b/rustfs/src/storage/rpc/mod.rs index 4b43fdbb8..ee3c51596 100644 --- a/rustfs/src/storage/rpc/mod.rs +++ b/rustfs/src/storage/rpc/mod.rs @@ -16,7 +16,10 @@ pub mod http_service; pub mod node_service; pub use http_service::InternodeRpcService; -pub use node_service::{HealControlRpcService, NodeService, make_heal_control_server, make_server}; +pub use node_service::{ + HealControlRpcService, NodeService, TierMutationControlRpcService, make_heal_control_server, make_server, + make_tier_mutation_control_server, +}; use rmp_serde::Serializer; use serde::Serialize; diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index da706455d..513745ef7 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -16,6 +16,7 @@ use crate::admin::service::{ config::{reload_dynamic_config_runtime_state, reload_runtime_config_snapshot}, site_replication::reload_site_replication_runtime_state, }; +use crate::storage::storage_api::ecstore_tier::tier_mutation_peer::{self, TierMutationPeerState as EcTierMutationPeerState}; #[cfg(test)] use crate::storage::storage_api::rpc_consumer::node_service::STORAGE_CLASS_SUB_SYS; #[cfg(test)] @@ -54,6 +55,7 @@ use tokio_stream::wrappers::ReceiverStream; use tokio_util::sync::CancellationToken; use tonic::{Request, Response, Status, Streaming}; use tracing::{debug, error, info, warn}; +use uuid::Uuid; pub(crate) mod heal; @@ -68,6 +70,10 @@ const EVENT_RPC_RESPONSE_EMITTED: &str = "rpc_response_emitted"; const EVENT_RPC_BACKGROUND_TASK_SPAWNED: &str = "rpc_background_task_spawned"; const EVENT_RPC_BACKGROUND_TASK_FAILED: &str = "rpc_background_task_failed"; const HEAL_CONTROL_REPLAY_CACHE_MAX_ENTRIES: usize = 4096; +const TIER_MUTATION_PEER_STATE_UNSPECIFIED_WIRE: i32 = 0; +const TIER_MUTATION_PEER_STATE_PREPARED_WIRE: i32 = 1; +const TIER_MUTATION_PEER_STATE_COMMITTED_WIRE: i32 = 2; +const TIER_MUTATION_PEER_STATE_ABORTED_WIRE: i32 = 3; #[derive(Debug)] struct HealControlReplayEntry { @@ -484,6 +490,203 @@ async fn execute_heal_control_envelope_with_manager( Ok(result) } +#[derive(Clone, Default)] +pub struct TierMutationControlRpcService { + context: Option>, +} + +impl std::fmt::Debug for TierMutationControlRpcService { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("TierMutationControlRpcService") + .field("context_present", &self.context.is_some()) + .finish() + } +} + +pub fn make_tier_mutation_control_server() -> TierMutationControlRpcService { + TierMutationControlRpcService { + context: runtime_sources::current_app_context(), + } +} + +#[cfg(test)] +pub(crate) fn make_tier_mutation_control_server_for_context( + context: Option>, +) -> TierMutationControlRpcService { + TierMutationControlRpcService { context } +} + +impl TierMutationControlRpcService { + fn resolve_object_store(&self) -> Option> { + let context = self.context.clone().or_else(runtime_sources::current_app_context); + runtime_sources::current_object_store_handle_for_context(context.as_deref()) + } + + async fn execute_tier_mutation( + &self, + request: &Request<()>, + version: u32, + phase: rustfs_protos::TierMutationRpcPhase, + mutation_id: &str, + canonical_payload: &Bytes, + ) -> Result, Status> { + validate_tier_mutation_payload_size(phase, canonical_payload.len())?; + let mutation_id = parse_tier_mutation_id(mutation_id)?; + let body = rustfs_protos::canonical_tier_mutation_rpc_body(version, phase, mutation_id, canonical_payload) + .map_err(|_| Status::invalid_argument("tier mutation request length cannot be represented"))?; + verify_tonic_canonical_body_digest(request, &body) + .map_err(|err| Status::permission_denied(format!("tier mutation authentication failed: {err}")))?; + let store = self + .resolve_object_store() + .ok_or_else(|| Status::failed_precondition("tier mutation object store is not initialized"))?; + + match tier_mutation_peer::handle_tier_mutation_peer_request(store, version, phase, mutation_id, canonical_payload).await { + Ok(outcome) => tier_mutation_control_response(TierMutationControlResponseInput { + version, + phase, + mutation_id, + canonical_payload, + success: true, + state: tier_mutation_peer_state_to_proto_wire(outcome.state), + applied: outcome.applied, + error_info: None, + }), + Err(err) => tier_mutation_control_response(TierMutationControlResponseInput { + version, + phase, + mutation_id, + canonical_payload, + success: false, + state: TIER_MUTATION_PEER_STATE_UNSPECIFIED_WIRE, + applied: false, + error_info: Some(err.to_string()), + }), + } + } +} + +#[tonic::async_trait] +impl tier_mutation_control_service_server::TierMutationControlService for TierMutationControlRpcService { + async fn prepare_tier_mutation( + &self, + request: Request, + ) -> Result, Status> { + let (metadata, extensions, inner) = request.into_parts(); + let request = Request::from_parts(metadata, extensions, ()); + self.execute_tier_mutation( + &request, + inner.version, + rustfs_protos::TierMutationRpcPhase::Prepare, + &inner.mutation_id, + &inner.canonical_payload, + ) + .await + } + + async fn commit_tier_mutation( + &self, + request: Request, + ) -> Result, Status> { + let (metadata, extensions, inner) = request.into_parts(); + let request = Request::from_parts(metadata, extensions, ()); + self.execute_tier_mutation( + &request, + inner.version, + rustfs_protos::TierMutationRpcPhase::Commit, + &inner.mutation_id, + &inner.canonical_payload, + ) + .await + } + + async fn abort_tier_mutation( + &self, + request: Request, + ) -> Result, Status> { + let (metadata, extensions, inner) = request.into_parts(); + let request = Request::from_parts(metadata, extensions, ()); + self.execute_tier_mutation( + &request, + inner.version, + rustfs_protos::TierMutationRpcPhase::Abort, + &inner.mutation_id, + &inner.canonical_payload, + ) + .await + } +} + +fn parse_tier_mutation_id(mutation_id: &str) -> Result { + let parsed = Uuid::parse_str(mutation_id).map_err(|_| Status::invalid_argument("tier mutation id is invalid"))?; + if parsed.to_string() != mutation_id { + return Err(Status::invalid_argument("tier mutation id is not canonical")); + } + Ok(parsed) +} + +fn validate_tier_mutation_payload_size(phase: rustfs_protos::TierMutationRpcPhase, payload_len: usize) -> Result<(), Status> { + let limit = match phase { + rustfs_protos::TierMutationRpcPhase::Prepare => rustfs_protos::TIER_MUTATION_RPC_MAX_PREPARE_PAYLOAD_SIZE, + rustfs_protos::TierMutationRpcPhase::Commit => rustfs_protos::TIER_MUTATION_RPC_MAX_COMMIT_PAYLOAD_SIZE, + rustfs_protos::TierMutationRpcPhase::Abort => { + if payload_len != 0 { + return Err(Status::invalid_argument("tier mutation abort payload must be empty")); + } + return Ok(()); + } + _ => return Err(Status::invalid_argument("tier mutation rpc phase is unsupported")), + }; + if payload_len > limit { + return Err(Status::invalid_argument("tier mutation payload exceeds size limit")); + } + Ok(()) +} + +struct TierMutationControlResponseInput<'a> { + version: u32, + phase: rustfs_protos::TierMutationRpcPhase, + mutation_id: Uuid, + canonical_payload: &'a [u8], + success: bool, + state: i32, + applied: bool, + error_info: Option, +} + +fn tier_mutation_control_response( + input: TierMutationControlResponseInput<'_>, +) -> Result, Status> { + let canonical_response = + rustfs_protos::canonical_tier_mutation_rpc_response_body(rustfs_protos::TierMutationRpcResponseProofInput { + version: input.version, + phase: input.phase, + mutation_id: input.mutation_id, + canonical_payload: input.canonical_payload, + success: input.success, + state: input.state, + applied: input.applied, + error_info: input.error_info.as_deref(), + }) + .map_err(|_| Status::internal("tier mutation response length cannot be represented"))?; + let response_proof = sign_tonic_rpc_response_proof(&canonical_response) + .map_err(|_| Status::internal("tier mutation response proof is unavailable"))?; + Ok(Response::new(TierMutationControlResponse { + success: input.success, + state: input.state, + applied: input.applied, + error_info: input.error_info, + response_proof: response_proof.into(), + })) +} + +fn tier_mutation_peer_state_to_proto_wire(state: EcTierMutationPeerState) -> i32 { + match state { + EcTierMutationPeerState::Prepared => TIER_MUTATION_PEER_STATE_PREPARED_WIRE, + EcTierMutationPeerState::Committed => TIER_MUTATION_PEER_STATE_COMMITTED_WIRE, + EcTierMutationPeerState::Aborted => TIER_MUTATION_PEER_STATE_ABORTED_WIRE, + } +} + #[tonic::async_trait] impl heal_control_service_server::HealControlService for HealControlRpcService { async fn heal_control(&self, request: Request) -> Result, Status> { @@ -1618,7 +1821,8 @@ mod tests { PEER_RESTSUB_SYS, SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC, STORAGE_CLASS_SUB_SYS, admit_heal_control_replay, background_rebalance_start_error_message, execute_heal_control_envelope_with_manager, initialize_heal_topology_fingerprint, make_heal_control_server, make_heal_control_server_with_cache, make_server, - remove_heal_control_replay, scanner_activity_response, stop_rebalance_response, + make_tier_mutation_control_server_for_context, remove_heal_control_replay, scanner_activity_response, + stop_rebalance_response, }; use crate::storage::rpc::node_service::heal::heal_topology_fingerprint; use crate::storage::storage_api::rpc_consumer::node_service::{HealBucketInfo, HealEndpoint}; @@ -1643,12 +1847,14 @@ mod tests { MakeVolumesRequest, Mss, PingRequest, ReadAllRequest, ReadAtRequest, ReadMultipleRequest, ReadVersionRequest, ReadXlRequest, ReloadPoolMetaRequest, ReloadSiteReplicationConfigRequest, RenameDataRequest, RenameFileRequest, RenamePartRequest, ScannerActivityRequest, ServerInfoRequest, SignalServiceRequest, StartProfilingRequest, - StatVolumeRequest, StopRebalanceRequest, UpdateMetacacheListingRequest, UpdateMetadataRequest, VerifyFileRequest, - WriteAllRequest, WriteMetadataRequest, WriteRequest, + StatVolumeRequest, StopRebalanceRequest, TierMutationPeerState, TierMutationPrepareRequest, + UpdateMetacacheListingRequest, UpdateMetadataRequest, VerifyFileRequest, WriteAllRequest, WriteMetadataRequest, + WriteRequest, heal_control_service_client::HealControlServiceClient, heal_control_service_server::{HealControlService as _, HealControlServiceServer}, node_service_client::NodeServiceClient, node_service_server::NodeServiceServer, + tier_mutation_control_service_server::TierMutationControlService as _, }; use std::{collections::HashMap, sync::Arc}; use time::OffsetDateTime; @@ -1967,6 +2173,24 @@ mod tests { .insert("x-rustfs-rpc-auth-version", "2".parse().expect("valid metadata value")); } + fn signed_tier_prepare_request(mutation_id: uuid::Uuid, canonical_payload: Bytes) -> Request { + let mut request = Request::new(TierMutationPrepareRequest { + version: rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION, + mutation_id: mutation_id.to_string(), + canonical_payload, + }); + let body = rustfs_protos::canonical_tier_mutation_rpc_body( + request.get_ref().version, + rustfs_protos::TierMutationRpcPhase::Prepare, + mutation_id, + &request.get_ref().canonical_payload, + ) + .expect("small request should encode"); + set_tonic_canonical_body_digest(&mut request, &body).expect("digest metadata should encode"); + mark_v2_authenticated(&mut request); + request + } + #[tokio::test] async fn heal_control_requires_body_bound_auth_before_topology_validation() { let service = make_heal_control_server(); @@ -2004,6 +2228,134 @@ mod tests { assert_eq!(unavailable.code(), tonic::Code::FailedPrecondition); } + #[tokio::test] + async fn tier_mutation_control_requires_body_bound_auth_before_store_lookup() { + let service = make_tier_mutation_control_server_for_context(None); + let mutation_id = uuid::Uuid::new_v4(); + let unsigned = service + .prepare_tier_mutation(Request::new(TierMutationPrepareRequest { + version: rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION, + mutation_id: mutation_id.to_string(), + canonical_payload: Bytes::from_static(b"intent"), + })) + .await + .expect_err("unsigned request must fail before store lookup"); + assert_eq!(unsigned.code(), tonic::Code::PermissionDenied); + + let mut tampered = signed_tier_prepare_request(mutation_id, Bytes::from_static(b"intent")); + let other_body = rustfs_protos::canonical_tier_mutation_rpc_body( + rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION, + rustfs_protos::TierMutationRpcPhase::Commit, + mutation_id, + b"intent", + ) + .expect("small request should encode"); + set_tonic_canonical_body_digest(&mut tampered, &other_body).expect("digest metadata should encode"); + let tampered = service + .prepare_tier_mutation(tampered) + .await + .expect_err("phase replay must fail body-bound authentication"); + assert_eq!(tampered.code(), tonic::Code::PermissionDenied); + + let signed = signed_tier_prepare_request(mutation_id, Bytes::from_static(b"intent")); + let unavailable = service + .prepare_tier_mutation(signed) + .await + .expect_err("authenticated request still requires initialized object store"); + assert_eq!(unavailable.code(), tonic::Code::FailedPrecondition); + } + + #[tokio::test] + async fn tier_mutation_control_requires_canonical_mutation_id() { + let service = make_tier_mutation_control_server_for_context(None); + let mutation_id = uuid::Uuid::new_v4().to_string().to_uppercase(); + let error = service + .prepare_tier_mutation(Request::new(TierMutationPrepareRequest { + version: rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION, + mutation_id, + canonical_payload: Bytes::from_static(b"intent"), + })) + .await + .expect_err("uppercase UUID must not pass canonical request binding"); + assert_eq!(error.code(), tonic::Code::InvalidArgument); + } + + #[tokio::test] + async fn tier_mutation_control_rejects_oversized_prepare_before_auth_and_store_lookup() { + let service = make_tier_mutation_control_server_for_context(None); + let mutation_id = uuid::Uuid::new_v4(); + let oversized = Bytes::from(vec![0; rustfs_protos::TIER_MUTATION_RPC_MAX_PREPARE_PAYLOAD_SIZE + 1]); + let error = service + .prepare_tier_mutation(Request::new(TierMutationPrepareRequest { + version: rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION, + mutation_id: mutation_id.to_string(), + canonical_payload: oversized, + })) + .await + .expect_err("oversized prepare must fail before digest construction"); + assert_eq!(error.code(), tonic::Code::InvalidArgument); + } + + #[test] + fn tier_mutation_peer_state_wire_constants_match_generated_proto() { + assert_eq!( + super::TIER_MUTATION_PEER_STATE_UNSPECIFIED_WIRE, + TierMutationPeerState::Unspecified as i32 + ); + assert_eq!(super::TIER_MUTATION_PEER_STATE_PREPARED_WIRE, TierMutationPeerState::Prepared as i32); + assert_eq!(super::TIER_MUTATION_PEER_STATE_COMMITTED_WIRE, TierMutationPeerState::Committed as i32); + assert_eq!(super::TIER_MUTATION_PEER_STATE_ABORTED_WIRE, TierMutationPeerState::Aborted as i32); + } + + #[test] + fn tier_mutation_control_response_proof_binds_request_and_result() { + let _ = rustfs_credentials::set_global_rpc_secret("tier-mutation-control-response-proof-test-secret".to_string()); + let mutation_id = uuid::Uuid::new_v4(); + let payload = b"canonical-intent-record"; + let response = super::tier_mutation_control_response(super::TierMutationControlResponseInput { + version: rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION, + phase: rustfs_protos::TierMutationRpcPhase::Prepare, + mutation_id, + canonical_payload: payload, + success: false, + state: TierMutationPeerState::Unspecified as i32, + applied: false, + error_info: Some("store failed".to_string()), + }) + .expect("response proof should be signed") + .into_inner(); + let canonical = + rustfs_protos::canonical_tier_mutation_rpc_response_body(rustfs_protos::TierMutationRpcResponseProofInput { + version: rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION, + phase: rustfs_protos::TierMutationRpcPhase::Prepare, + mutation_id, + canonical_payload: payload, + success: false, + state: TierMutationPeerState::Unspecified as i32, + applied: false, + error_info: Some("store failed"), + }) + .expect("small mutation response should encode"); + crate::storage::storage_api::verify_tonic_rpc_response_proof(&canonical, &response.response_proof) + .expect("proof must authenticate the exact response"); + + let tampered = + rustfs_protos::canonical_tier_mutation_rpc_response_body(rustfs_protos::TierMutationRpcResponseProofInput { + version: rustfs_protos::TIER_MUTATION_RPC_PROTOCOL_VERSION, + phase: rustfs_protos::TierMutationRpcPhase::Prepare, + mutation_id, + canonical_payload: payload, + success: true, + state: TierMutationPeerState::Unspecified as i32, + applied: false, + error_info: Some("store failed"), + }) + .expect("small mutation response should encode"); + let error = crate::storage::storage_api::verify_tonic_rpc_response_proof(&tampered, &response.response_proof) + .expect_err("proof must reject a tampered success flag"); + assert_eq!(error.to_string(), "Invalid RPC response proof"); + } + #[tokio::test] async fn heal_control_rejects_oversized_command_before_canonical_copy() { let service = make_heal_control_server(); diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index dd09bfb7f..2e16c52b2 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -333,7 +333,9 @@ pub(crate) mod timeout_wrapper_consumer { pub(crate) mod tonic_service_consumer { #[cfg(test)] pub(crate) use super::super::tonic_service::{heal_topology_fingerprint, make_heal_control_server_for_source}; - pub(crate) use super::super::tonic_service::{make_heal_control_server_with_cache, make_server}; + pub(crate) use super::super::tonic_service::{ + make_heal_control_server_with_cache, make_server, make_tier_mutation_control_server, + }; } #[cfg(test)] @@ -512,7 +514,7 @@ pub(crate) mod ecstore_storage { pub(crate) mod ecstore_tier { pub(crate) use rustfs_ecstore::api::tier::tier::{TierConfigMgr, TierConfigUpdateError}; - pub(crate) use rustfs_ecstore::api::tier::{tier, tier_admin, tier_config, tier_handlers}; + pub(crate) use rustfs_ecstore::api::tier::{tier, tier_admin, tier_config, tier_handlers, tier_mutation_peer}; // Shared lifecycle/tier test utilities behind ecstore's `test-util` feature // (rustfs/backlog#1148 ilm-6). Only linked into test builds. #[cfg(test)] diff --git a/rustfs/src/storage/tonic_service.rs b/rustfs/src/storage/tonic_service.rs index 3d2921b3b..e56a8f235 100644 --- a/rustfs/src/storage/tonic_service.rs +++ b/rustfs/src/storage/tonic_service.rs @@ -15,6 +15,6 @@ pub(crate) use crate::storage::rpc::node_service::make_heal_control_server_with_cache; #[cfg(test)] pub(crate) use crate::storage::rpc::node_service::{heal::heal_topology_fingerprint, make_heal_control_server_for_source}; -pub use crate::storage::rpc::{make_heal_control_server, make_server}; +pub use crate::storage::rpc::{make_heal_control_server, make_server, make_tier_mutation_control_server}; #[allow(dead_code)] pub type NodeService = crate::storage::rpc::NodeService; diff --git a/rustfs/src/storage_api.rs b/rustfs/src/storage_api.rs index 6a5361362..b691899d1 100644 --- a/rustfs/src/storage_api.rs +++ b/rustfs/src/storage_api.rs @@ -135,7 +135,7 @@ pub(crate) mod server { heal_topology_fingerprint, make_heal_control_server_for_source, }; pub(crate) use crate::storage::storage_api::tonic_service_consumer::{ - make_heal_control_server_with_cache, make_server, + make_heal_control_server_with_cache, make_server, make_tier_mutation_control_server, }; } }