mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-16 18:08:21 +00:00
feat(heal): aggregate replacement recovery status (#5916)
Add a replacement recovery peer RPC so Admin v4 can distinguish definitive cluster proofs from unsupported, unavailable, or conflicting peer state without extending the existing background heal v3/v1 status protocol. Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -22,7 +22,7 @@ use crate::server::ADMIN_PREFIX;
|
||||
use crate::server::RemoteAddr;
|
||||
use crate::storage::rpc::node_service::heal::{
|
||||
HealControlCoordinator, NodeHealProgress, NodeHealStatusSnapshot, capture_node_heal_status, decode_node_heal_status,
|
||||
heal_control_coordinator, heal_topology_fingerprint,
|
||||
decode_node_replacement_recovery_status, heal_control_coordinator, heal_topology_fingerprint,
|
||||
};
|
||||
use bytes::Bytes;
|
||||
use futures_util::future::join_all;
|
||||
@@ -40,7 +40,7 @@ use rustfs_utils::path::path_join;
|
||||
use s3s::header::{CONTENT_LENGTH, CONTENT_TYPE};
|
||||
use s3s::{Body, S3Request, S3Response, S3Result, s3_error};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::collections::HashSet;
|
||||
use std::collections::{BTreeMap, HashSet};
|
||||
use std::future::Future;
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
@@ -55,7 +55,7 @@ const EVENT_ADMIN_REQUEST_FAILED: &str = "admin_request_failed";
|
||||
const EVENT_ADMIN_RESPONSE_EMITTED: &str = "admin_response_emitted";
|
||||
const PEER_HEAL_STATUS_TIMEOUT: Duration = Duration::from_secs(5);
|
||||
pub(crate) const REPLACEMENT_RECOVERY_STATUS_ROUTE_SUFFIX: &str = "/v4/heal/replacement-recovery";
|
||||
const REPLACEMENT_RECOVERY_STATUS_CONTRACT_VERSION: u32 = 1;
|
||||
const REPLACEMENT_RECOVERY_STATUS_CONTRACT_VERSION: u32 = 2;
|
||||
|
||||
#[derive(Debug, Default, Serialize, Deserialize)]
|
||||
struct HealInitParams {
|
||||
@@ -529,6 +529,204 @@ async fn read_cluster_heal_status(
|
||||
merge_peer_heal_statuses(snapshots, peer_statuses, expected_nodes, topology_complete)
|
||||
}
|
||||
|
||||
async fn query_peer_replacement_recovery_status<E>(
|
||||
host: &str,
|
||||
request: impl std::future::Future<Output = Result<Option<Vec<u8>>, E>>,
|
||||
request_timeout: Duration,
|
||||
) -> Result<Option<rustfs_heal::ReplacementRecoverySnapshot>, String>
|
||||
where
|
||||
E: std::fmt::Display,
|
||||
{
|
||||
match timeout(request_timeout, request).await {
|
||||
Ok(Ok(Some(status))) => decode_node_replacement_recovery_status(&status)
|
||||
.map(|snapshot| Some(snapshot.snapshot))
|
||||
.map_err(|err| format!("peer {host}: {err}")),
|
||||
Ok(Ok(None)) => Ok(None),
|
||||
Ok(Err(err)) => Err(format!("peer {host}: {err}")),
|
||||
Err(_) => Err(format!("peer {host}: replacement recovery status timed out")),
|
||||
}
|
||||
}
|
||||
|
||||
fn canonical_replacement_records(
|
||||
records: &[rustfs_heal::ReplacementRecoveryRecord],
|
||||
) -> Vec<rustfs_heal::ReplacementRecoveryRecord> {
|
||||
let mut records = records.to_vec();
|
||||
for record in &mut records {
|
||||
record.target_slots.sort();
|
||||
}
|
||||
records.sort_by(|left, right| {
|
||||
(
|
||||
&left.task_id,
|
||||
&left.generation,
|
||||
&left.set_disk_id,
|
||||
left.state,
|
||||
&left.target_slots,
|
||||
&left.verified_at,
|
||||
&left.reason,
|
||||
)
|
||||
.cmp(&(
|
||||
&right.task_id,
|
||||
&right.generation,
|
||||
&right.set_disk_id,
|
||||
right.state,
|
||||
&right.target_slots,
|
||||
&right.verified_at,
|
||||
&right.reason,
|
||||
))
|
||||
});
|
||||
records
|
||||
}
|
||||
|
||||
fn merge_replacement_recovery_records(
|
||||
snapshots: &[rustfs_heal::ReplacementRecoverySnapshot],
|
||||
) -> Vec<rustfs_heal::ReplacementRecoveryRecord> {
|
||||
let mut records = BTreeMap::<String, rustfs_heal::ReplacementRecoveryRecord>::new();
|
||||
for snapshot in snapshots {
|
||||
for record in canonical_replacement_records(&snapshot.records) {
|
||||
match records.entry(record.task_id.clone()) {
|
||||
std::collections::btree_map::Entry::Vacant(entry) => {
|
||||
entry.insert(record);
|
||||
}
|
||||
std::collections::btree_map::Entry::Occupied(entry) if entry.get() == &record => {}
|
||||
std::collections::btree_map::Entry::Occupied(mut entry) => {
|
||||
let task_id = entry.key().clone();
|
||||
entry.insert(rustfs_heal::ReplacementRecoveryRecord {
|
||||
task_id,
|
||||
state: rustfs_heal::ReplacementRecoveryState::Unknown,
|
||||
generation: None,
|
||||
set_disk_id: None,
|
||||
target_slots: Vec::new(),
|
||||
reason: Some("peer replacement recovery records disagree".to_string()),
|
||||
verified_at: None,
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
records.into_values().collect()
|
||||
}
|
||||
|
||||
fn aggregate_replacement_recovery_cluster_status(
|
||||
mut snapshots: Vec<rustfs_heal::ReplacementRecoverySnapshot>,
|
||||
peer_statuses: Vec<Result<Option<rustfs_heal::ReplacementRecoverySnapshot>, String>>,
|
||||
expected_nodes: usize,
|
||||
topology_complete: bool,
|
||||
) -> ReplacementRecoveryClusterStatus {
|
||||
let mut reason = None::<&'static str>;
|
||||
if !topology_complete {
|
||||
reason = Some("peer_topology_incomplete");
|
||||
}
|
||||
|
||||
for status in peer_statuses {
|
||||
match status {
|
||||
Ok(Some(snapshot)) => snapshots.push(snapshot),
|
||||
Ok(None) => {
|
||||
reason.get_or_insert("peer_replacement_recovery_status_unsupported");
|
||||
}
|
||||
Err(_) => {
|
||||
reason.get_or_insert("peer_replacement_recovery_status_unavailable");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if snapshots.len() != expected_nodes {
|
||||
reason.get_or_insert("peer_replacement_recovery_status_incomplete");
|
||||
}
|
||||
if snapshots.iter().any(|snapshot| !snapshot.definitive) {
|
||||
reason.get_or_insert("peer_replacement_recovery_status_not_definitive");
|
||||
}
|
||||
|
||||
let first_records = snapshots
|
||||
.first()
|
||||
.map(|snapshot| canonical_replacement_records(&snapshot.records));
|
||||
if let Some(first_records) = &first_records
|
||||
&& snapshots
|
||||
.iter()
|
||||
.skip(1)
|
||||
.any(|snapshot| canonical_replacement_records(&snapshot.records) != *first_records)
|
||||
{
|
||||
reason.get_or_insert("peer_replacement_recovery_status_conflict");
|
||||
}
|
||||
|
||||
let records = merge_replacement_recovery_records(&snapshots);
|
||||
if records
|
||||
.iter()
|
||||
.any(|record| matches!(record.state, rustfs_heal::ReplacementRecoveryState::Unknown))
|
||||
{
|
||||
reason.get_or_insert("peer_replacement_recovery_status_conflict");
|
||||
}
|
||||
|
||||
ReplacementRecoveryClusterStatus {
|
||||
definitive: reason.is_none(),
|
||||
reason: reason.map(str::to_string),
|
||||
records,
|
||||
}
|
||||
}
|
||||
|
||||
async fn read_cluster_replacement_recovery_status(
|
||||
local: rustfs_heal::ReplacementRecoverySnapshot,
|
||||
notification_system: Option<&crate::admin::storage_api::runtime_sources::NotificationSys>,
|
||||
expected_nodes: usize,
|
||||
) -> S3Result<ReplacementRecoveryClusterStatus> {
|
||||
if expected_nodes == 1 {
|
||||
return Ok(aggregate_replacement_recovery_cluster_status(
|
||||
vec![local],
|
||||
Vec::new(),
|
||||
expected_nodes,
|
||||
true,
|
||||
));
|
||||
}
|
||||
let Some(notification_system) = notification_system else {
|
||||
return Ok(ReplacementRecoveryClusterStatus {
|
||||
definitive: false,
|
||||
reason: Some("notification_system_unavailable".to_string()),
|
||||
records: local.records,
|
||||
});
|
||||
};
|
||||
let topology_complete = peer_topology_complete(
|
||||
expected_nodes,
|
||||
notification_system.peer_clients.len(),
|
||||
notification_system
|
||||
.peer_clients
|
||||
.iter()
|
||||
.filter(|client| client.is_none())
|
||||
.count(),
|
||||
notification_system.all_peer_clients.len(),
|
||||
notification_system
|
||||
.all_peer_clients
|
||||
.iter()
|
||||
.filter(|client| client.is_some())
|
||||
.count(),
|
||||
);
|
||||
if !topology_complete {
|
||||
warn!(
|
||||
event = EVENT_ADMIN_REQUEST_FAILED,
|
||||
component = LOG_COMPONENT_ADMIN_API,
|
||||
subsystem = LOG_SUBSYSTEM_HEAL_ADMIN,
|
||||
operation = "replacement_recovery_status",
|
||||
result = "degraded",
|
||||
reason = "peer_topology_incomplete",
|
||||
"cluster replacement recovery status will be partial"
|
||||
);
|
||||
}
|
||||
|
||||
let peer_statuses = join_all(notification_system.peer_clients.iter().map(|client| async move {
|
||||
let Some(client) = client else {
|
||||
return Err("configured peer is unavailable".to_string());
|
||||
};
|
||||
let host = client.host.to_string();
|
||||
query_peer_replacement_recovery_status(&host, client.replacement_recovery_status(), PEER_HEAL_STATUS_TIMEOUT).await
|
||||
}))
|
||||
.await;
|
||||
|
||||
Ok(aggregate_replacement_recovery_cluster_status(
|
||||
vec![local],
|
||||
peer_statuses,
|
||||
expected_nodes,
|
||||
topology_complete,
|
||||
))
|
||||
}
|
||||
|
||||
fn cluster_heal_control_unavailable(reason: &str) -> s3s::S3Error {
|
||||
warn!(
|
||||
event = EVENT_ADMIN_REQUEST_FAILED,
|
||||
@@ -1036,21 +1234,20 @@ pub(crate) struct ReplacementRecoveryStatusResponse {
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub(crate) struct ReplacementRecoveryClusterStatus {
|
||||
pub definitive: bool,
|
||||
pub reason: &'static str,
|
||||
pub reason: Option<String>,
|
||||
pub records: Vec<rustfs_heal::ReplacementRecoveryRecord>,
|
||||
}
|
||||
|
||||
fn build_replacement_recovery_status_response(
|
||||
local: rustfs_heal::ReplacementRecoverySnapshot,
|
||||
cluster: ReplacementRecoveryClusterStatus,
|
||||
) -> ReplacementRecoveryStatusResponse {
|
||||
ReplacementRecoveryStatusResponse {
|
||||
contract_version: REPLACEMENT_RECOVERY_STATUS_CONTRACT_VERSION,
|
||||
path: format!("{}{}", ADMIN_PREFIX, REPLACEMENT_RECOVERY_STATUS_ROUTE_SUFFIX),
|
||||
scope: "local_survivor_disks",
|
||||
scope: "cluster_survivor_disks",
|
||||
local,
|
||||
cluster: ReplacementRecoveryClusterStatus {
|
||||
definitive: false,
|
||||
reason: "distributed replacement recovery status requires a peer capability RPC",
|
||||
},
|
||||
cluster,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1347,8 +1544,21 @@ impl Operation for ReplacementRecoveryStatusHandler {
|
||||
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
||||
validate_heal_admin_request(&req).await?;
|
||||
|
||||
let Some(context) = app_context_from_req(&req) else {
|
||||
return Err(cluster_heal_status_unavailable("server_not_initialized"));
|
||||
};
|
||||
let Some(endpoints) = context.endpoints().handle() else {
|
||||
return Err(cluster_heal_status_unavailable("endpoint_topology_unavailable"));
|
||||
};
|
||||
let expected_nodes = endpoints.get_nodes().len();
|
||||
if expected_nodes == 0 {
|
||||
return Err(cluster_heal_status_unavailable("endpoint_topology_empty"));
|
||||
}
|
||||
let notification_system = context.notification_system().handle();
|
||||
let local = rustfs_heal::current_replacement_recovery_snapshot().await;
|
||||
let response = build_replacement_recovery_status_response(local);
|
||||
let cluster =
|
||||
read_cluster_replacement_recovery_status(local.clone(), notification_system.as_deref(), expected_nodes).await?;
|
||||
let response = build_replacement_recovery_status_response(local, cluster);
|
||||
let body = encode_replacement_recovery_status(&response)?;
|
||||
info!(
|
||||
event = EVENT_ADMIN_RESPONSE_EMITTED,
|
||||
@@ -1368,14 +1578,17 @@ mod tests {
|
||||
use super::extract_heal_init_params;
|
||||
use super::{
|
||||
BackgroundHealProgress, HealInitParams, HealResp, HealRuntimeState, aggregate_cluster_heal_status,
|
||||
background_heal_runtime_state, build_heal_channel_request, build_replacement_recovery_status_response,
|
||||
encode_background_heal_status, encode_heal_start_success, encode_heal_task_status, execute_after_heal_control_capability,
|
||||
heal_channel_response_items, heal_channel_response_progress, heal_channel_response_summary, json_response,
|
||||
map_heal_response, map_root_heal_status, merge_peer_heal_statuses, peer_topology_complete, query_peer_heal_status,
|
||||
aggregate_replacement_recovery_cluster_status, background_heal_runtime_state, build_heal_channel_request,
|
||||
build_replacement_recovery_status_response, encode_background_heal_status, encode_heal_start_success,
|
||||
encode_heal_task_status, execute_after_heal_control_capability, heal_channel_response_items,
|
||||
heal_channel_response_progress, heal_channel_response_summary, json_response, map_heal_response, map_root_heal_status,
|
||||
merge_peer_heal_statuses, peer_topology_complete, query_peer_heal_status, query_peer_replacement_recovery_status,
|
||||
reject_heal_admission, should_handle_root_heal_directly, validate_heal_request_mode, validate_heal_target,
|
||||
};
|
||||
use crate::admin::storage_api::error::StorageError;
|
||||
use crate::storage::rpc::node_service::heal::{NodeHealProgress, NodeHealStatusSnapshot};
|
||||
use crate::storage::rpc::node_service::heal::{
|
||||
NodeHealProgress, NodeHealStatusSnapshot, NodeReplacementRecoveryStatusSnapshot, encode_node_replacement_recovery_status,
|
||||
};
|
||||
use bytes::Bytes;
|
||||
use http::StatusCode;
|
||||
use http::Uri;
|
||||
@@ -1394,6 +1607,26 @@ mod tests {
|
||||
use tokio::sync::mpsc;
|
||||
use tokio::time::Duration;
|
||||
|
||||
fn replacement_record(task_id: &str) -> rustfs_heal::ReplacementRecoveryRecord {
|
||||
rustfs_heal::ReplacementRecoveryRecord {
|
||||
task_id: task_id.to_string(),
|
||||
state: rustfs_heal::ReplacementRecoveryState::Completed,
|
||||
generation: Some(task_id.to_string()),
|
||||
set_disk_id: Some("pool_0_set_0".to_string()),
|
||||
target_slots: vec!["http://node-a:9000/mnt/disk1".to_string()],
|
||||
reason: None,
|
||||
verified_at: Some(42),
|
||||
}
|
||||
}
|
||||
|
||||
fn replacement_snapshot(task_id: &str) -> rustfs_heal::ReplacementRecoverySnapshot {
|
||||
rustfs_heal::ReplacementRecoverySnapshot {
|
||||
records: vec![replacement_record(task_id)],
|
||||
definitive: true,
|
||||
reason: None,
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn cluster_capability_gate_runs_before_execution() {
|
||||
let executed = AtomicBool::new(false);
|
||||
@@ -1421,32 +1654,94 @@ mod tests {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn replacement_recovery_status_response_never_claims_cluster_completion() {
|
||||
let response = build_replacement_recovery_status_response(rustfs_heal::ReplacementRecoverySnapshot {
|
||||
records: vec![rustfs_heal::ReplacementRecoveryRecord {
|
||||
task_id: "11111111-1111-4111-8111-111111111111".to_string(),
|
||||
state: rustfs_heal::ReplacementRecoveryState::Completed,
|
||||
generation: Some("11111111-1111-4111-8111-111111111111".to_string()),
|
||||
set_disk_id: Some("pool_0_set_0".to_string()),
|
||||
target_slots: vec!["http://node-a:9000/mnt/disk1".to_string()],
|
||||
reason: None,
|
||||
verified_at: Some(42),
|
||||
}],
|
||||
definitive: true,
|
||||
reason: None,
|
||||
});
|
||||
fn replacement_recovery_status_response_reports_cluster_proof() {
|
||||
let local = replacement_snapshot("11111111-1111-4111-8111-111111111111");
|
||||
let cluster = aggregate_replacement_recovery_cluster_status(vec![local.clone()], Vec::new(), 1, true);
|
||||
let response = build_replacement_recovery_status_response(local, cluster);
|
||||
let value = serde_json::to_value(response).expect("replacement status response should serialize");
|
||||
|
||||
assert_eq!(value["contractVersion"], 1);
|
||||
assert_eq!(value["contractVersion"], 2);
|
||||
assert_eq!(value["path"], "/rustfs/admin/v4/heal/replacement-recovery");
|
||||
assert_eq!(value["scope"], "local_survivor_disks");
|
||||
assert_eq!(value["scope"], "cluster_survivor_disks");
|
||||
assert_eq!(value["local"]["definitive"], true);
|
||||
assert_eq!(value["local"]["records"][0]["state"], "completed");
|
||||
assert_eq!(value["cluster"]["definitive"], false);
|
||||
assert_eq!(
|
||||
value["cluster"]["reason"],
|
||||
"distributed replacement recovery status requires a peer capability RPC"
|
||||
);
|
||||
assert_eq!(value["cluster"]["definitive"], true);
|
||||
assert!(value["cluster"]["reason"].is_null());
|
||||
assert_eq!(value["cluster"]["records"][0]["state"], "completed");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn replacement_recovery_cluster_requires_peer_capability() {
|
||||
let local = replacement_snapshot("11111111-1111-4111-8111-111111111111");
|
||||
let cluster = aggregate_replacement_recovery_cluster_status(vec![local], vec![Ok(None)], 2, true);
|
||||
|
||||
assert!(!cluster.definitive);
|
||||
assert_eq!(cluster.reason.as_deref(), Some("peer_replacement_recovery_status_unsupported"));
|
||||
assert_eq!(cluster.records[0].state, rustfs_heal::ReplacementRecoveryState::Completed);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn replacement_recovery_cluster_requires_matching_peer_records() {
|
||||
let local = replacement_snapshot("11111111-1111-4111-8111-111111111111");
|
||||
let peer = replacement_snapshot("22222222-2222-4222-8222-222222222222");
|
||||
let cluster = aggregate_replacement_recovery_cluster_status(vec![local], vec![Ok(Some(peer))], 2, true);
|
||||
|
||||
assert!(!cluster.definitive);
|
||||
assert_eq!(cluster.reason.as_deref(), Some("peer_replacement_recovery_status_conflict"));
|
||||
assert_eq!(cluster.records.len(), 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn replacement_recovery_cluster_rejects_conflicting_records_inside_a_snapshot() {
|
||||
let mut local = replacement_snapshot("11111111-1111-4111-8111-111111111111");
|
||||
let mut conflicting = replacement_record("11111111-1111-4111-8111-111111111111");
|
||||
conflicting.set_disk_id = Some("pool_0_set_1".to_string());
|
||||
local.records.push(conflicting);
|
||||
let cluster = aggregate_replacement_recovery_cluster_status(vec![local], Vec::new(), 1, true);
|
||||
|
||||
assert!(!cluster.definitive);
|
||||
assert_eq!(cluster.reason.as_deref(), Some("peer_replacement_recovery_status_conflict"));
|
||||
assert_eq!(cluster.records[0].state, rustfs_heal::ReplacementRecoveryState::Unknown);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn replacement_recovery_cluster_is_definitive_when_peers_match() {
|
||||
let local = replacement_snapshot("11111111-1111-4111-8111-111111111111");
|
||||
let peer = replacement_snapshot("11111111-1111-4111-8111-111111111111");
|
||||
let cluster = aggregate_replacement_recovery_cluster_status(vec![local], vec![Ok(Some(peer))], 2, true);
|
||||
|
||||
assert!(cluster.definitive);
|
||||
assert!(cluster.reason.is_none());
|
||||
assert_eq!(cluster.records.len(), 1);
|
||||
assert_eq!(cluster.records[0].state, rustfs_heal::ReplacementRecoveryState::Completed);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn peer_replacement_recovery_status_decodes_supported_peers() {
|
||||
let payload = encode_node_replacement_recovery_status(&NodeReplacementRecoveryStatusSnapshot::new(replacement_snapshot(
|
||||
"11111111-1111-4111-8111-111111111111",
|
||||
)))
|
||||
.expect("peer status should encode");
|
||||
let decoded =
|
||||
query_peer_replacement_recovery_status("node-a", async { Ok::<_, String>(Some(payload)) }, Duration::from_secs(1))
|
||||
.await
|
||||
.expect("peer query should succeed")
|
||||
.expect("peer should support status");
|
||||
|
||||
assert_eq!(decoded.records[0].state, rustfs_heal::ReplacementRecoveryState::Completed);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn peer_replacement_recovery_status_treats_missing_rpc_as_unknown() {
|
||||
let decoded = query_peer_replacement_recovery_status(
|
||||
"node-a",
|
||||
async { Ok::<Option<Vec<u8>>, String>(None) },
|
||||
Duration::from_secs(1),
|
||||
)
|
||||
.await
|
||||
.expect("missing rpc should not be a transport failure");
|
||||
|
||||
assert!(decoded.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -1934,6 +1934,32 @@ impl Node for NodeService {
|
||||
}
|
||||
}
|
||||
|
||||
async fn replacement_recovery_status(
|
||||
&self,
|
||||
_request: Request<ReplacementRecoveryStatusRequest>,
|
||||
) -> Result<Response<ReplacementRecoveryStatusResponse>, Status> {
|
||||
if self.resolve_object_store().is_none() {
|
||||
return Ok(Response::new(ReplacementRecoveryStatusResponse {
|
||||
success: false,
|
||||
recovery_status: Bytes::new(),
|
||||
error_info: Some("storage layer not initialized".to_string()),
|
||||
}));
|
||||
}
|
||||
let snapshot = heal::capture_node_replacement_recovery_status().await;
|
||||
match heal::encode_node_replacement_recovery_status(&snapshot) {
|
||||
Ok(recovery_status) => Ok(Response::new(ReplacementRecoveryStatusResponse {
|
||||
success: true,
|
||||
recovery_status: recovery_status.into(),
|
||||
error_info: None,
|
||||
})),
|
||||
Err(err) => Ok(Response::new(ReplacementRecoveryStatusResponse {
|
||||
success: false,
|
||||
recovery_status: Bytes::new(),
|
||||
error_info: Some(err),
|
||||
})),
|
||||
}
|
||||
}
|
||||
|
||||
async fn get_metacache_listing(
|
||||
&self,
|
||||
_request: Request<GetMetacacheListingRequest>,
|
||||
|
||||
@@ -27,6 +27,8 @@ use super::super::encode_msgpack_map;
|
||||
|
||||
const NODE_HEAL_STATUS_VERSION: u8 = 1;
|
||||
const NODE_HEAL_STATUS_MAX_SIZE: usize = 64 * 1024;
|
||||
const NODE_REPLACEMENT_RECOVERY_STATUS_VERSION: u8 = 1;
|
||||
const NODE_REPLACEMENT_RECOVERY_STATUS_MAX_SIZE: usize = 64 * 1024;
|
||||
const HEAL_TOPOLOGY_FINGERPRINT_DOMAIN: &[u8] = b"rustfs-heal-topology-v1\0";
|
||||
|
||||
fn chrono_to_jiff_timestamp(timestamp: chrono::DateTime<chrono::Utc>) -> Timestamp {
|
||||
@@ -286,18 +288,115 @@ pub(crate) fn decode_node_heal_status(data: &[u8]) -> Result<NodeHealStatusSnaps
|
||||
Ok(snapshot)
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
|
||||
#[serde(rename_all = "camelCase", deny_unknown_fields)]
|
||||
pub(crate) struct NodeReplacementRecoveryStatusSnapshot {
|
||||
version: u8,
|
||||
pub snapshot: rustfs_heal::ReplacementRecoverySnapshot,
|
||||
}
|
||||
|
||||
impl NodeReplacementRecoveryStatusSnapshot {
|
||||
pub(crate) fn new(snapshot: rustfs_heal::ReplacementRecoverySnapshot) -> Self {
|
||||
Self {
|
||||
version: NODE_REPLACEMENT_RECOVERY_STATUS_VERSION,
|
||||
snapshot,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
#[serde(rename_all = "camelCase", deny_unknown_fields)]
|
||||
struct NodeReplacementRecoveryStatusSnapshotWire {
|
||||
version: u8,
|
||||
snapshot: ReplacementRecoverySnapshotWire,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
#[serde(rename_all = "camelCase", deny_unknown_fields)]
|
||||
struct ReplacementRecoverySnapshotWire {
|
||||
records: Vec<ReplacementRecoveryRecordWire>,
|
||||
definitive: bool,
|
||||
reason: Option<String>,
|
||||
}
|
||||
|
||||
impl From<ReplacementRecoverySnapshotWire> for rustfs_heal::ReplacementRecoverySnapshot {
|
||||
fn from(snapshot: ReplacementRecoverySnapshotWire) -> Self {
|
||||
Self {
|
||||
records: snapshot.records.into_iter().map(Into::into).collect(),
|
||||
definitive: snapshot.definitive,
|
||||
reason: snapshot.reason,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
#[serde(rename_all = "camelCase", deny_unknown_fields)]
|
||||
struct ReplacementRecoveryRecordWire {
|
||||
task_id: String,
|
||||
state: rustfs_heal::ReplacementRecoveryState,
|
||||
generation: Option<String>,
|
||||
set_disk_id: Option<String>,
|
||||
target_slots: Vec<String>,
|
||||
reason: Option<String>,
|
||||
verified_at: Option<u64>,
|
||||
}
|
||||
|
||||
impl From<ReplacementRecoveryRecordWire> for rustfs_heal::ReplacementRecoveryRecord {
|
||||
fn from(record: ReplacementRecoveryRecordWire) -> Self {
|
||||
Self {
|
||||
task_id: record.task_id,
|
||||
state: record.state,
|
||||
generation: record.generation,
|
||||
set_disk_id: record.set_disk_id,
|
||||
target_slots: record.target_slots,
|
||||
reason: record.reason,
|
||||
verified_at: record.verified_at,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) async fn capture_node_replacement_recovery_status() -> NodeReplacementRecoveryStatusSnapshot {
|
||||
NodeReplacementRecoveryStatusSnapshot::new(rustfs_heal::current_replacement_recovery_snapshot().await)
|
||||
}
|
||||
|
||||
pub(crate) fn encode_node_replacement_recovery_status(
|
||||
snapshot: &NodeReplacementRecoveryStatusSnapshot,
|
||||
) -> Result<Vec<u8>, String> {
|
||||
encode_msgpack_map(snapshot).map_err(|err| format!("failed to encode replacement recovery status: {err}"))
|
||||
}
|
||||
|
||||
pub(crate) fn decode_node_replacement_recovery_status(data: &[u8]) -> Result<NodeReplacementRecoveryStatusSnapshot, String> {
|
||||
if data.len() > NODE_REPLACEMENT_RECOVERY_STATUS_MAX_SIZE {
|
||||
return Err("replacement recovery status exceeds size limit".to_string());
|
||||
}
|
||||
let mut deserializer = Deserializer::new(Cursor::new(data));
|
||||
let wire = NodeReplacementRecoveryStatusSnapshotWire::deserialize(&mut deserializer)
|
||||
.map_err(|err| format!("failed to decode replacement recovery status: {err}"))?;
|
||||
if usize::try_from(deserializer.get_ref().position()).ok() != Some(data.len()) {
|
||||
return Err("replacement recovery status contains trailing data".to_string());
|
||||
}
|
||||
if wire.version != NODE_REPLACEMENT_RECOVERY_STATUS_VERSION {
|
||||
return Err(format!("unsupported replacement recovery status version: {}", wire.version));
|
||||
}
|
||||
Ok(NodeReplacementRecoveryStatusSnapshot {
|
||||
version: wire.version,
|
||||
snapshot: wire.snapshot.into(),
|
||||
})
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::{
|
||||
NODE_HEAL_STATUS_MAX_SIZE, NODE_HEAL_STATUS_VERSION, NodeHealProgress, NodeHealStatusSnapshot, decode_node_heal_status,
|
||||
encode_node_heal_status, heal_control_coordinator, heal_topology_fingerprint,
|
||||
NODE_HEAL_STATUS_MAX_SIZE, NODE_HEAL_STATUS_VERSION, NodeHealProgress, NodeHealStatusSnapshot,
|
||||
NodeReplacementRecoveryStatusSnapshot, decode_node_heal_status, decode_node_replacement_recovery_status,
|
||||
encode_node_heal_status, encode_node_replacement_recovery_status, heal_control_coordinator, heal_topology_fingerprint,
|
||||
};
|
||||
use crate::storage::storage_api::{
|
||||
Endpoint,
|
||||
ecstore_layout::{EndpointServerPools, Endpoints, PoolEndpoints},
|
||||
};
|
||||
use chrono::SecondsFormat;
|
||||
use rustfs_heal::HealOperationsSnapshot;
|
||||
use rustfs_heal::{HealOperationsSnapshot, ReplacementRecoveryRecord, ReplacementRecoverySnapshot, ReplacementRecoveryState};
|
||||
use rustfs_scanner::scanner::BackgroundHealInfo;
|
||||
|
||||
fn topology_endpoints(last_host: &str) -> EndpointServerPools {
|
||||
@@ -566,4 +665,72 @@ mod tests {
|
||||
.contains("size limit")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn node_replacement_recovery_status_round_trip_is_versioned() {
|
||||
let snapshot = NodeReplacementRecoveryStatusSnapshot::new(ReplacementRecoverySnapshot {
|
||||
records: vec![ReplacementRecoveryRecord {
|
||||
task_id: "11111111-1111-4111-8111-111111111111".to_string(),
|
||||
state: ReplacementRecoveryState::Completed,
|
||||
generation: Some("11111111-1111-4111-8111-111111111111".to_string()),
|
||||
set_disk_id: Some("pool_0_set_0".to_string()),
|
||||
target_slots: vec!["http://node-a:9000/mnt/disk1".to_string()],
|
||||
reason: None,
|
||||
verified_at: Some(42),
|
||||
}],
|
||||
definitive: true,
|
||||
reason: None,
|
||||
});
|
||||
|
||||
let encoded = encode_node_replacement_recovery_status(&snapshot).expect("snapshot should encode");
|
||||
let decoded = decode_node_replacement_recovery_status(&encoded).expect("snapshot should decode");
|
||||
|
||||
assert_eq!(decoded, snapshot);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn node_replacement_recovery_status_rejects_unknown_trailing_and_oversized_data() {
|
||||
let mut snapshot = NodeReplacementRecoveryStatusSnapshot::new(ReplacementRecoverySnapshot {
|
||||
records: Vec::new(),
|
||||
definitive: true,
|
||||
reason: None,
|
||||
});
|
||||
snapshot.version += 1;
|
||||
let encoded = encode_node_replacement_recovery_status(&snapshot).expect("snapshot should encode");
|
||||
assert!(
|
||||
decode_node_replacement_recovery_status(&encoded)
|
||||
.expect_err("unknown version should fail closed")
|
||||
.contains("unsupported replacement recovery status version")
|
||||
);
|
||||
|
||||
snapshot.version = super::NODE_REPLACEMENT_RECOVERY_STATUS_VERSION;
|
||||
let mut encoded = encode_node_replacement_recovery_status(&snapshot).expect("snapshot should encode");
|
||||
encoded.push(0);
|
||||
assert!(
|
||||
decode_node_replacement_recovery_status(&encoded)
|
||||
.expect_err("trailing data must fail")
|
||||
.contains("trailing")
|
||||
);
|
||||
assert!(
|
||||
decode_node_replacement_recovery_status(&vec![0; super::NODE_REPLACEMENT_RECOVERY_STATUS_MAX_SIZE + 1])
|
||||
.expect_err("oversized status must fail")
|
||||
.contains("size limit")
|
||||
);
|
||||
|
||||
let fixture = serde_json::json!({
|
||||
"version": super::NODE_REPLACEMENT_RECOVERY_STATUS_VERSION,
|
||||
"snapshot": {
|
||||
"records": [],
|
||||
"definitive": true,
|
||||
"reason": null,
|
||||
"futureField": true,
|
||||
}
|
||||
});
|
||||
let encoded = rmp_serde::to_vec_named(&fixture).expect("fixture should encode");
|
||||
assert!(
|
||||
decode_node_replacement_recovery_status(&encoded)
|
||||
.expect_err("nested unknown field must fail")
|
||||
.contains("futureField")
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user