fix(storage): sync transition tier config across peers (#2918)

Co-authored-by: houseme <housemecn@gmail.com>
This commit is contained in:
cxymds
2026-05-11 23:23:29 +08:00
committed by GitHub
parent d6813b53a2
commit 8ef2c0a5e1
3 changed files with 182 additions and 84 deletions
+14
View File
@@ -920,4 +920,18 @@ mod tests {
assert!(msg.contains("1 failure(s)"));
assert!(msg.contains("peer[0]"));
}
#[tokio::test]
async fn load_transition_tier_config_reports_unreachable_peers() {
let sys = NotificationSys {
peer_clients: vec![None],
all_peer_clients: Vec::new(),
};
let results = sys.load_transition_tier_config().await;
assert_eq!(results.len(), 1);
assert!(results[0].host.is_empty());
assert!(results[0].err.is_some());
assert!(results[0].err.as_ref().unwrap().to_string().contains("peer is not reachable"));
}
}
+134 -77
View File
@@ -18,7 +18,7 @@ use crate::{
auth::validate_admin_request,
router::{AdminOperation, Operation, S3Router},
},
app::context::resolve_tier_config_handle,
app::context::{resolve_object_store_handle, resolve_tier_config_handle},
auth::{check_key_valid, get_session_token},
server::{ADMIN_PREFIX, RemoteAddr},
};
@@ -34,6 +34,7 @@ use rustfs_ecstore::{
bucket::lifecycle::tier_last_day_stats::DailyAllTierStats,
client::admin_handler_utils::AdminError,
config::storageclass,
notification_sys::get_global_notification_sys,
tier::{
tier::{ERR_TIER_BACKEND_IN_USE, ERR_TIER_BACKEND_NOT_EMPTY, ERR_TIER_MISSING_CREDENTIALS},
tier_admin::TierCreds,
@@ -53,6 +54,7 @@ use s3s::{
use serde_urlencoded::from_bytes;
use std::collections::HashMap;
use time::OffsetDateTime;
use tokio::spawn;
use tracing::{debug, warn};
#[derive(Debug, Clone, serde::Deserialize, Default)]
@@ -83,6 +85,22 @@ pub struct AddTierQuery {
pub struct AddTier {}
fn spawn_transition_tier_config_propagation(action: &'static str) {
if let Some(notification_sys) = get_global_notification_sys() {
spawn(async move {
for peer_result in notification_sys.load_transition_tier_config().await {
if let Some(err) = peer_result.err {
warn!(
host = if peer_result.host.is_empty() { "<unknown>" } else { peer_result.host.as_str() },
error = %err,
"tier {action} propagation failed after local save"
);
}
}
});
}
}
fn resolve_tier_name(uri: &Uri, params: &Params<'_, '_>) -> S3Result<String> {
if let Some(tier) = params.get("tier") {
let decoded = percent_decode_str(tier)
@@ -240,44 +258,57 @@ impl Operation for AddTier {
&_ => (),
}
let tier_config_mgr_handle = resolve_tier_config_handle();
let mut tier_config_mgr = tier_config_mgr_handle.write().await;
//tier_config_mgr.reload(api);
if let Err(err) = tier_config_mgr.add(args, force).await {
return if err.code == ERR_TIER_ALREADY_EXISTS.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierNameAlreadyExist".into()),
"tier name already exists!",
))
} else if err.code == ERR_TIER_NAME_NOT_UPPERCASE.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierNameNotUppercase".into()),
"tier name not uppercase!",
))
} else if err.code == ERR_TIER_BACKEND_IN_USE.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierNameBackendInUse!".into()),
"tier name backend in use!",
))
} else if err.code == ERR_TIER_CONNECT_ERR.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierConnectError".into()),
"tier connect error!",
))
} else if err.code == ERR_TIER_INVALID_CREDENTIALS.code {
Err(S3Error::with_message(S3ErrorCode::Custom(err.code.clone().into()), err.message))
} else {
warn!("tier_config_mgr add failed, e: {:?}", err);
Err(S3Error::with_message(
let Some(store) = resolve_object_store_handle() else {
return Err(s3_error!(InvalidRequest, "object store not init"));
};
{
let tier_config_mgr_handle = resolve_tier_config_handle();
let mut tier_config_mgr = tier_config_mgr_handle.write().await;
if let Err(err) = tier_config_mgr.reload(store).await {
warn!("tier_config_mgr reload failed, e: {:?}", err);
return Err(S3Error::with_message(
S3ErrorCode::Custom("TierAddFailed".into()),
format!("tier add failed. {err}"),
))
};
}
if let Err(e) = tier_config_mgr.save().await {
warn!("tier_config_mgr save failed, e: {:?}", e);
return Err(S3Error::with_message(S3ErrorCode::Custom("TierAddFailed".into()), "tier save failed"));
format!("tier reload failed. {err}"),
));
}
if let Err(err) = tier_config_mgr.add(args, force).await {
return if err.code == ERR_TIER_ALREADY_EXISTS.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierNameAlreadyExist".into()),
"tier name already exists!",
))
} else if err.code == ERR_TIER_NAME_NOT_UPPERCASE.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierNameNotUppercase".into()),
"tier name not uppercase!",
))
} else if err.code == ERR_TIER_BACKEND_IN_USE.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierNameBackendInUse!".into()),
"tier name backend in use!",
))
} else if err.code == ERR_TIER_CONNECT_ERR.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierConnectError".into()),
"tier connect error!",
))
} else if err.code == ERR_TIER_INVALID_CREDENTIALS.code {
Err(S3Error::with_message(S3ErrorCode::Custom(err.code.clone().into()), err.message))
} else {
warn!("tier_config_mgr add failed, e: {:?}", err);
Err(S3Error::with_message(
S3ErrorCode::Custom("TierAddFailed".into()),
format!("tier add failed. {err}"),
))
};
}
if let Err(e) = tier_config_mgr.save().await {
warn!("tier_config_mgr save failed, e: {:?}", e);
return Err(S3Error::with_message(S3ErrorCode::Custom("TierAddFailed".into()), "tier save failed"));
}
}
spawn_transition_tier_config_propagation("add");
let mut header = HeaderMap::new();
header.insert(CONTENT_TYPE, "application/json".parse().unwrap());
@@ -333,29 +364,42 @@ impl Operation for EditTier {
let tier_name = params.get("tiername").map(|s| s.to_string()).unwrap_or_default();
let tier_config_mgr_handle = resolve_tier_config_handle();
let mut tier_config_mgr = tier_config_mgr_handle.write().await;
//tier_config_mgr.reload(api);
if let Err(err) = tier_config_mgr.edit(&tier_name, creds).await {
return if err.code == ERR_TIER_NOT_FOUND.code {
Err(S3Error::with_message(S3ErrorCode::Custom("TierNotFound".into()), "tier not found!"))
} else if err.code == ERR_TIER_MISSING_CREDENTIALS.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierMissingCredentials".into()),
"tier missing credentials!",
))
} else {
warn!("tier_config_mgr edit failed, e: {:?}", err);
Err(S3Error::with_message(
let Some(store) = resolve_object_store_handle() else {
return Err(s3_error!(InvalidRequest, "object store not init"));
};
{
let tier_config_mgr_handle = resolve_tier_config_handle();
let mut tier_config_mgr = tier_config_mgr_handle.write().await;
if let Err(err) = tier_config_mgr.reload(store).await {
warn!("tier_config_mgr reload failed, e: {:?}", err);
return Err(S3Error::with_message(
S3ErrorCode::Custom("TierEditFailed".into()),
format!("tier edit failed. {err}"),
))
};
}
if let Err(e) = tier_config_mgr.save().await {
warn!("tier_config_mgr save failed, e: {:?}", e);
return Err(S3Error::with_message(S3ErrorCode::Custom("TierEditFailed".into()), "tier save failed"));
format!("tier reload failed. {err}"),
));
}
if let Err(err) = tier_config_mgr.edit(&tier_name, creds).await {
return if err.code == ERR_TIER_NOT_FOUND.code {
Err(S3Error::with_message(S3ErrorCode::Custom("TierNotFound".into()), "tier not found!"))
} else if err.code == ERR_TIER_MISSING_CREDENTIALS.code {
Err(S3Error::with_message(
S3ErrorCode::Custom("TierMissingCredentials".into()),
"tier missing credentials!",
))
} else {
warn!("tier_config_mgr edit failed, e: {:?}", err);
Err(S3Error::with_message(
S3ErrorCode::Custom("TierEditFailed".into()),
format!("tier edit failed. {err}"),
))
};
}
if let Err(e) = tier_config_mgr.save().await {
warn!("tier_config_mgr save failed, e: {:?}", e);
return Err(S3Error::with_message(S3ErrorCode::Custom("TierEditFailed".into()), "tier save failed"));
}
}
spawn_transition_tier_config_propagation("edit");
let mut header = HeaderMap::new();
header.insert(CONTENT_TYPE, "application/json".parse().unwrap());
@@ -457,27 +501,40 @@ impl Operation for RemoveTier {
let tier_name = params.get("tiername").map(|s| s.to_string()).unwrap_or_default();
let tier_config_mgr_handle = resolve_tier_config_handle();
let mut tier_config_mgr = tier_config_mgr_handle.write().await;
//tier_config_mgr.reload(api);
if let Err(err) = tier_config_mgr.remove(&tier_name, force).await {
return if err.code == ERR_TIER_NOT_FOUND.code {
Err(S3Error::with_message(S3ErrorCode::Custom("TierNotFound".into()), "tier not found."))
} else if err.code == ERR_TIER_BACKEND_NOT_EMPTY.code {
Err(S3Error::with_message(S3ErrorCode::Custom("TierNameBackendInUse".into()), "tier is used."))
} else {
warn!("tier_config_mgr remove failed, e: {:?}", err);
Err(S3Error::with_message(
S3ErrorCode::Custom("TierRemoveFailed".into()),
format!("tier remove failed. {err}"),
))
};
}
let Some(store) = resolve_object_store_handle() else {
return Err(s3_error!(InvalidRequest, "object store not init"));
};
if let Err(e) = tier_config_mgr.save().await {
warn!("tier_config_mgr save failed, e: {:?}", e);
return Err(S3Error::with_message(S3ErrorCode::Custom("TierRemoveFailed".into()), "tier save failed"));
{
let tier_config_mgr_handle = resolve_tier_config_handle();
let mut tier_config_mgr = tier_config_mgr_handle.write().await;
if let Err(err) = tier_config_mgr.reload(store).await {
warn!("tier_config_mgr reload failed, e: {:?}", err);
return Err(S3Error::with_message(
S3ErrorCode::Custom("TierRemoveFailed".into()),
format!("tier reload failed. {err}"),
));
}
if let Err(err) = tier_config_mgr.remove(&tier_name, force).await {
return if err.code == ERR_TIER_NOT_FOUND.code {
Err(S3Error::with_message(S3ErrorCode::Custom("TierNotFound".into()), "tier not found."))
} else if err.code == ERR_TIER_BACKEND_NOT_EMPTY.code {
Err(S3Error::with_message(S3ErrorCode::Custom("TierNameBackendInUse".into()), "tier is used."))
} else {
warn!("tier_config_mgr remove failed, e: {:?}", err);
Err(S3Error::with_message(
S3ErrorCode::Custom("TierRemoveFailed".into()),
format!("tier remove failed. {err}"),
))
};
}
if let Err(e) = tier_config_mgr.save().await {
warn!("tier_config_mgr save failed, e: {:?}", e);
return Err(S3Error::with_message(S3ErrorCode::Custom("TierRemoveFailed".into()), "tier save failed"));
}
}
spawn_transition_tier_config_propagation("remove");
let mut header = HeaderMap::new();
header.insert(CONTENT_TYPE, "application/json".parse().unwrap());
+34 -7
View File
@@ -26,6 +26,7 @@ use rustfs_ecstore::{
UpdateMetadataOpts, error::DiskError,
},
get_global_lock_client,
global::GLOBAL_TierConfigMgr,
metrics_realtime::{CollectMetricsOpts, MetricType, collect_local_metrics},
new_object_layer_fn,
rpc::{LocalPeerS3Client, PeerS3Client},
@@ -894,7 +895,23 @@ impl Node for NodeService {
&self,
_request: Request<LoadTransitionTierConfigRequest>,
) -> Result<Response<LoadTransitionTierConfigResponse>, Status> {
Err(unimplemented_rpc("load_transition_tier_config"))
let Some(store) = new_object_layer_fn() else {
return Ok(Response::new(LoadTransitionTierConfigResponse {
success: false,
error_info: Some("errServerNotInitialized".to_string()),
}));
};
match GLOBAL_TierConfigMgr.write().await.reload(store).await {
Ok(_) => Ok(Response::new(LoadTransitionTierConfigResponse {
success: true,
error_info: None,
})),
Err(err) => Ok(Response::new(LoadTransitionTierConfigResponse {
success: false,
error_info: Some(err.to_string()),
})),
}
}
}
@@ -2228,6 +2245,22 @@ mod tests {
assert!(load_response.error_info.unwrap().contains("errServerNotInitialized"));
}
#[tokio::test]
#[ignore = "requires isolated global object layer state"]
async fn test_load_transition_tier_config_no_object_layer() {
let service = create_test_node_service();
let response = service
.load_transition_tier_config(Request::new(LoadTransitionTierConfigRequest::default()))
.await;
assert!(response.is_ok());
let load_response = response.unwrap().into_inner();
assert!(!load_response.success);
assert!(load_response.error_info.is_some());
assert!(load_response.error_info.unwrap().contains("errServerNotInitialized"));
}
#[tokio::test]
async fn test_delete_bucket_metadata() {
let service = create_test_node_service();
@@ -2461,12 +2494,6 @@ mod tests {
.await,
"update_metacache_listing",
);
assert_unimplemented_status(
service
.load_transition_tier_config(Request::new(LoadTransitionTierConfigRequest::default()))
.await,
"load_transition_tier_config",
);
}
async fn connect_test_node_service_client() -> NodeServiceClient<tonic::transport::Channel> {