mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-22 12:26:37 +00:00
Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 8df19b7b8b |
@@ -585,12 +585,9 @@ impl VersionsHistogram {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Replication statistics for a single target.
|
/// Replication statistics for a single target
|
||||||
///
|
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
|
||||||
/// Renamed from `ReplicationStats`; serde field names are preserved
|
pub struct ReplicationStats {
|
||||||
/// byte-identically to maintain wire compatibility with existing snapshots.
|
|
||||||
#[derive(Debug, Default, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
|
||||||
pub struct ReplicationTargetUsage {
|
|
||||||
pub pending_size: u64,
|
pub pending_size: u64,
|
||||||
pub replicated_size: u64,
|
pub replicated_size: u64,
|
||||||
pub failed_size: u64,
|
pub failed_size: u64,
|
||||||
@@ -603,7 +600,7 @@ pub struct ReplicationTargetUsage {
|
|||||||
pub replicated_count: u64,
|
pub replicated_count: u64,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl ReplicationTargetUsage {
|
impl ReplicationStats {
|
||||||
pub fn is_empty(&self) -> bool {
|
pub fn is_empty(&self) -> bool {
|
||||||
let Self {
|
let Self {
|
||||||
pending_size,
|
pending_size,
|
||||||
@@ -639,7 +636,7 @@ impl ReplicationTargetUsage {
|
|||||||
/// Replication statistics for all targets
|
/// Replication statistics for all targets
|
||||||
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
|
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
|
||||||
pub struct ReplicationAllStats {
|
pub struct ReplicationAllStats {
|
||||||
pub targets: HashMap<String, ReplicationTargetUsage>,
|
pub targets: HashMap<String, ReplicationStats>,
|
||||||
pub replica_size: u64,
|
pub replica_size: u64,
|
||||||
pub replica_count: u64,
|
pub replica_count: u64,
|
||||||
}
|
}
|
||||||
@@ -652,7 +649,7 @@ impl ReplicationAllStats {
|
|||||||
targets,
|
targets,
|
||||||
} = self;
|
} = self;
|
||||||
|
|
||||||
*replica_size == 0 && *replica_count == 0 && targets.values().all(ReplicationTargetUsage::is_empty)
|
*replica_size == 0 && *replica_count == 0 && targets.values().all(ReplicationStats::is_empty)
|
||||||
}
|
}
|
||||||
|
|
||||||
#[deprecated(note = "use is_empty instead")]
|
#[deprecated(note = "use is_empty instead")]
|
||||||
@@ -2469,7 +2466,7 @@ mod tests {
|
|||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn replication_stats_empty_checks_every_field() {
|
fn replication_stats_empty_checks_every_field() {
|
||||||
type SetField = fn(&mut ReplicationTargetUsage);
|
type SetField = fn(&mut ReplicationStats);
|
||||||
|
|
||||||
let cases: [(&str, SetField); 10] = [
|
let cases: [(&str, SetField); 10] = [
|
||||||
("pending_size", |stats| stats.pending_size = 1),
|
("pending_size", |stats| stats.pending_size = 1),
|
||||||
@@ -2484,9 +2481,9 @@ mod tests {
|
|||||||
("replicated_count", |stats| stats.replicated_count = 1),
|
("replicated_count", |stats| stats.replicated_count = 1),
|
||||||
];
|
];
|
||||||
|
|
||||||
assert!(ReplicationTargetUsage::default().is_empty());
|
assert!(ReplicationStats::default().is_empty());
|
||||||
for (field, set_nonzero) in cases {
|
for (field, set_nonzero) in cases {
|
||||||
let mut stats = ReplicationTargetUsage::default();
|
let mut stats = ReplicationStats::default();
|
||||||
set_nonzero(&mut stats);
|
set_nonzero(&mut stats);
|
||||||
assert!(!stats.is_empty(), "{field} must make replication stats non-empty");
|
assert!(!stats.is_empty(), "{field} must make replication stats non-empty");
|
||||||
}
|
}
|
||||||
@@ -2517,17 +2514,17 @@ mod tests {
|
|||||||
}
|
}
|
||||||
|
|
||||||
let empty_targets = ReplicationAllStats {
|
let empty_targets = ReplicationAllStats {
|
||||||
targets: HashMap::from([("arn:test:empty".to_string(), ReplicationTargetUsage::default())]),
|
targets: HashMap::from([("arn:test:empty".to_string(), ReplicationStats::default())]),
|
||||||
..Default::default()
|
..Default::default()
|
||||||
};
|
};
|
||||||
assert!(empty_targets.is_empty(), "all-empty targets must keep aggregate stats empty");
|
assert!(empty_targets.is_empty(), "all-empty targets must keep aggregate stats empty");
|
||||||
|
|
||||||
let stats = ReplicationAllStats {
|
let stats = ReplicationAllStats {
|
||||||
targets: HashMap::from([
|
targets: HashMap::from([
|
||||||
("arn:test:empty".to_string(), ReplicationTargetUsage::default()),
|
("arn:test:empty".to_string(), ReplicationStats::default()),
|
||||||
(
|
(
|
||||||
"arn:test:non-empty".to_string(),
|
"arn:test:non-empty".to_string(),
|
||||||
ReplicationTargetUsage {
|
ReplicationStats {
|
||||||
pending_count: 1,
|
pending_count: 1,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
@@ -2568,7 +2565,7 @@ mod tests {
|
|||||||
replication_stats: Some(ReplicationAllStats {
|
replication_stats: Some(ReplicationAllStats {
|
||||||
targets: HashMap::from([(
|
targets: HashMap::from([(
|
||||||
"arn:test:pending".to_string(),
|
"arn:test:pending".to_string(),
|
||||||
ReplicationTargetUsage {
|
ReplicationStats {
|
||||||
pending_count: 1,
|
pending_count: 1,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
@@ -2717,7 +2714,7 @@ mod tests {
|
|||||||
targets: HashMap::from([
|
targets: HashMap::from([
|
||||||
(
|
(
|
||||||
"arn:self-only".to_string(),
|
"arn:self-only".to_string(),
|
||||||
ReplicationTargetUsage {
|
ReplicationStats {
|
||||||
pending_size: 7,
|
pending_size: 7,
|
||||||
pending_count: 1,
|
pending_count: 1,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
@@ -2725,7 +2722,7 @@ mod tests {
|
|||||||
),
|
),
|
||||||
(
|
(
|
||||||
"arn:shared".to_string(),
|
"arn:shared".to_string(),
|
||||||
ReplicationTargetUsage {
|
ReplicationStats {
|
||||||
failed_size: 3,
|
failed_size: 3,
|
||||||
failed_count: 1,
|
failed_count: 1,
|
||||||
missed_threshold_size: 2,
|
missed_threshold_size: 2,
|
||||||
@@ -2744,7 +2741,7 @@ mod tests {
|
|||||||
targets: HashMap::from([
|
targets: HashMap::from([
|
||||||
(
|
(
|
||||||
"arn:shared".to_string(),
|
"arn:shared".to_string(),
|
||||||
ReplicationTargetUsage {
|
ReplicationStats {
|
||||||
failed_size: 5,
|
failed_size: 5,
|
||||||
failed_count: 2,
|
failed_count: 2,
|
||||||
after_threshold_size: 4,
|
after_threshold_size: 4,
|
||||||
@@ -2754,7 +2751,7 @@ mod tests {
|
|||||||
),
|
),
|
||||||
(
|
(
|
||||||
"arn:other-only".to_string(),
|
"arn:other-only".to_string(),
|
||||||
ReplicationTargetUsage {
|
ReplicationStats {
|
||||||
replicated_size: 11,
|
replicated_size: 11,
|
||||||
replicated_count: 3,
|
replicated_count: 3,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
@@ -2996,9 +2993,7 @@ mod tests {
|
|||||||
fn replication_target_deserialization_preserves_large_historical_maps() {
|
fn replication_target_deserialization_preserves_large_historical_maps() {
|
||||||
let mut stats = ReplicationAllStats::default();
|
let mut stats = ReplicationAllStats::default();
|
||||||
for index in 0..=1024 {
|
for index in 0..=1024 {
|
||||||
stats
|
stats.targets.insert(format!("target-{index}"), ReplicationStats::default());
|
||||||
.targets
|
|
||||||
.insert(format!("target-{index}"), ReplicationTargetUsage::default());
|
|
||||||
}
|
}
|
||||||
let encoded = rmp_serde::to_vec_named(&stats).expect("large replication target fixture should encode");
|
let encoded = rmp_serde::to_vec_named(&stats).expect("large replication target fixture should encode");
|
||||||
let decoded = rmp_serde::from_slice::<ReplicationAllStats>(&encoded)
|
let decoded = rmp_serde::from_slice::<ReplicationAllStats>(&encoded)
|
||||||
@@ -3007,47 +3002,6 @@ mod tests {
|
|||||||
assert_eq!(decoded.targets.len(), stats.targets.len());
|
assert_eq!(decoded.targets.len(), stats.targets.len());
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Round-trip test: encoding a [`ReplicationTargetUsage`] and decoding it back
|
|
||||||
/// must produce the exact same value. This guards against accidental serde
|
|
||||||
/// field-name drift during the `ReplicationStats` -> `ReplicationTargetUsage`
|
|
||||||
/// rename. Wire-level field names are the serialized Rust field identifiers,
|
|
||||||
/// which must remain byte-identical.
|
|
||||||
#[test]
|
|
||||||
fn replication_target_usage_rmp_round_trip() {
|
|
||||||
let original = ReplicationTargetUsage {
|
|
||||||
pending_size: 100,
|
|
||||||
replicated_size: 2_000,
|
|
||||||
failed_size: 50,
|
|
||||||
failed_count: 3,
|
|
||||||
pending_count: 7,
|
|
||||||
missed_threshold_size: 11,
|
|
||||||
after_threshold_size: 22,
|
|
||||||
missed_threshold_count: 1,
|
|
||||||
after_threshold_count: 2,
|
|
||||||
replicated_count: 99,
|
|
||||||
};
|
|
||||||
|
|
||||||
let buf = rmp_serde::to_vec_named(&original).expect("encode ReplicationTargetUsage to msgpack");
|
|
||||||
let decoded: ReplicationTargetUsage = rmp_serde::from_slice(&buf).expect("decode ReplicationTargetUsage from msgpack");
|
|
||||||
assert_eq!(original, decoded, "round-trip through rmp must preserve every field");
|
|
||||||
|
|
||||||
// Also verify that encoding as an unnamed sequence and then decoding
|
|
||||||
// with named fields produces the correct mapping (this catches reordering).
|
|
||||||
let named_buf = rmp_serde::to_vec_named(&original).expect("re-encode for field-name pinning");
|
|
||||||
// Spot-check that known field names appear in the named encoding.
|
|
||||||
let named_str = String::from_utf8_lossy(&named_buf);
|
|
||||||
assert!(named_str.contains("pending_size"), "field 'pending_size' must survive the rename");
|
|
||||||
assert!(named_str.contains("replicated_size"), "field 'replicated_size' must survive the rename");
|
|
||||||
assert!(
|
|
||||||
named_str.contains("missed_threshold_size"),
|
|
||||||
"field 'missed_threshold_size' must survive the rename"
|
|
||||||
);
|
|
||||||
assert!(
|
|
||||||
named_str.contains("after_threshold_count"),
|
|
||||||
"field 'after_threshold_count' must survive the rename"
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn checked_merge_rejects_noncanonical_histograms_without_mutation() {
|
fn checked_merge_rejects_noncanonical_histograms_without_mutation() {
|
||||||
let mut entry = DataUsageEntry {
|
let mut entry = DataUsageEntry {
|
||||||
|
|||||||
@@ -380,10 +380,24 @@ mod tests {
|
|||||||
cluster.start_node(1).await?;
|
cluster.start_node(1).await?;
|
||||||
|
|
||||||
let status_url = format!("{}/rustfs/admin/v3/background-heal/status", cluster.nodes[0].url);
|
let status_url = format!("{}/rustfs/admin/v3/background-heal/status", cluster.nodes[0].url);
|
||||||
let status_body = signed_admin_post(&status_url, None, &cluster.access_key, &cluster.secret_key).await?;
|
let mut recovered = serde_json::Value::Null;
|
||||||
assert!(
|
for _ in 0..60 {
|
||||||
!status_body.contains("MissingContentLength"),
|
let status_body = signed_admin_post(&status_url, None, &cluster.access_key, &cluster.secret_key).await?;
|
||||||
"background heal status should not fail without an explicit Content-Length: {status_body}"
|
assert!(
|
||||||
|
!status_body.contains("MissingContentLength"),
|
||||||
|
"background heal status should not fail without an explicit Content-Length: {status_body}"
|
||||||
|
);
|
||||||
|
recovered = serde_json::from_str(&status_body)
|
||||||
|
.map_err(|err| format!("background heal status is not JSON ({err}): {status_body}"))?;
|
||||||
|
if recovered["clusterStatusComplete"] == serde_json::Value::Bool(true) {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
sleep(Duration::from_secs(1)).await;
|
||||||
|
}
|
||||||
|
assert_eq!(
|
||||||
|
recovered["clusterStatusComplete"],
|
||||||
|
serde_json::Value::Bool(true),
|
||||||
|
"cluster heal status should recover before root heal starts: {recovered}"
|
||||||
);
|
);
|
||||||
|
|
||||||
let heal_body = r#"{"recursive":true,"dryRun":false,"remove":false,"recreate":true,"scanMode":2,"updateParity":false,"nolock":false}"#;
|
let heal_body = r#"{"recursive":true,"dryRun":false,"remove":false,"recreate":true,"scanMode":2,"updateParity":false,"nolock":false}"#;
|
||||||
|
|||||||
@@ -50,10 +50,10 @@ use rustfs_protos::evict_failed_connection;
|
|||||||
use rustfs_protos::proto_gen::node_service::RenamePartRequest;
|
use rustfs_protos::proto_gen::node_service::RenamePartRequest;
|
||||||
use rustfs_protos::proto_gen::node_service::{
|
use rustfs_protos::proto_gen::node_service::{
|
||||||
BatchReadVersionRequest, BatchReadVersionResponse, CheckPartsRequest, DeletePathsRequest, DeleteRequest,
|
BatchReadVersionRequest, BatchReadVersionResponse, CheckPartsRequest, DeletePathsRequest, DeleteRequest,
|
||||||
DeleteVersionRequest, DeleteVersionsRequest, DeleteVersionsResponse, DeleteVolumeRequest, DiskInfoRequest, ListDirRequest,
|
DeleteVersionRequest, DeleteVersionsRequest, DeleteVolumeRequest, DiskInfoRequest, ListDirRequest, ListVolumesRequest,
|
||||||
ListVolumesRequest, MakeVolumeRequest, MakeVolumesRequest, PreparePartTransactionRequest, ReadAllRequest,
|
MakeVolumeRequest, MakeVolumesRequest, PreparePartTransactionRequest, ReadAllRequest, ReadMetadataRequest,
|
||||||
ReadMetadataRequest, ReadMultipleRequest, ReadMultipleResponse, ReadPartsRequest, ReadVersionRequest, ReadXlRequest,
|
ReadMultipleRequest, ReadMultipleResponse, ReadPartsRequest, ReadVersionRequest, ReadXlRequest, RenameDataRequest,
|
||||||
RenameDataRequest, RenameFileRequest, SettlePartTransactionRequest, SnapshotLeaseReleaseRequest, SnapshotLeaseRenewRequest,
|
RenameFileRequest, SettlePartTransactionRequest, SnapshotLeaseReleaseRequest, SnapshotLeaseRenewRequest,
|
||||||
SnapshotLeaseRequest, SnapshotLeaseResponse, StatVolumeRequest, UpdateMetadataRequest, VerifyFileRequest, WriteAllRequest,
|
SnapshotLeaseRequest, SnapshotLeaseResponse, StatVolumeRequest, UpdateMetadataRequest, VerifyFileRequest, WriteAllRequest,
|
||||||
WriteMetadataRequest, node_service_client::NodeServiceClient,
|
WriteMetadataRequest, node_service_client::NodeServiceClient,
|
||||||
};
|
};
|
||||||
@@ -112,28 +112,6 @@ const EVENT_REMOTE_DISK_RPC: &str = "remote_disk_rpc";
|
|||||||
const SNAPSHOT_LEASE_PROTOCOL_VERSION: u32 = 1;
|
const SNAPSHOT_LEASE_PROTOCOL_VERSION: u32 = 1;
|
||||||
pub const REMOTE_SNAPSHOT_LEASE_TTL: Duration = Duration::from_secs(60);
|
pub const REMOTE_SNAPSHOT_LEASE_TTL: Duration = Duration::from_secs(60);
|
||||||
|
|
||||||
fn decode_delete_versions_errors(response: DeleteVersionsResponse, expected_len: usize) -> Vec<Option<Error>> {
|
|
||||||
if !response.item_errors.is_empty() {
|
|
||||||
if response.item_errors.len() != expected_len {
|
|
||||||
return vec![Some(Error::other("malformed delete_versions item errors")); expected_len];
|
|
||||||
}
|
|
||||||
return response
|
|
||||||
.item_errors
|
|
||||||
.into_iter()
|
|
||||||
.map(|error| (error.code != 0).then(|| error.into()))
|
|
||||||
.collect();
|
|
||||||
}
|
|
||||||
|
|
||||||
if response.errors.len() != expected_len {
|
|
||||||
return vec![Some(Error::other("malformed delete_versions errors")); expected_len];
|
|
||||||
}
|
|
||||||
response
|
|
||||||
.errors
|
|
||||||
.into_iter()
|
|
||||||
.map(|error| (!error.is_empty()).then(|| Error::other(error)))
|
|
||||||
.collect()
|
|
||||||
}
|
|
||||||
|
|
||||||
fn snapshot_lease_token_from_response(response: SnapshotLeaseResponse) -> Result<SnapshotLeaseToken> {
|
fn snapshot_lease_token_from_response(response: SnapshotLeaseResponse) -> Result<SnapshotLeaseToken> {
|
||||||
if !response.success {
|
if !response.success {
|
||||||
return Err(response.error.unwrap_or_default().into());
|
return Err(response.error.unwrap_or_default().into());
|
||||||
@@ -2428,6 +2406,8 @@ impl DiskAPI for RemoteDisk {
|
|||||||
return errors;
|
return errors;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TODO(backlog): replace string errors with typed `StorageError` variants
|
||||||
|
|
||||||
let result = self
|
let result = self
|
||||||
.execute_with_timeout(
|
.execute_with_timeout(
|
||||||
|| async {
|
|| async {
|
||||||
@@ -2459,7 +2439,17 @@ impl DiskAPI for RemoteDisk {
|
|||||||
}
|
}
|
||||||
return errors;
|
return errors;
|
||||||
}
|
}
|
||||||
decode_delete_versions_errors(response, versions.len())
|
response
|
||||||
|
.errors
|
||||||
|
.iter()
|
||||||
|
.map(|error| {
|
||||||
|
if error.is_empty() {
|
||||||
|
None
|
||||||
|
} else {
|
||||||
|
Some(Error::other(error.to_string()))
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.collect()
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tracing::instrument(level = "trace", skip_all)]
|
#[tracing::instrument(level = "trace", skip_all)]
|
||||||
@@ -3770,63 +3760,6 @@ mod tests {
|
|||||||
|
|
||||||
static INIT: Once = Once::new();
|
static INIT: Once = Once::new();
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn delete_versions_response_preserves_typed_item_errors() {
|
|
||||||
let errors = decode_delete_versions_errors(
|
|
||||||
DeleteVersionsResponse {
|
|
||||||
success: true,
|
|
||||||
errors: vec!["file not found".to_string(), String::new()],
|
|
||||||
error: None,
|
|
||||||
item_errors: vec![
|
|
||||||
rustfs_protos::proto_gen::node_service::Error {
|
|
||||||
code: DiskError::FileNotFound.to_u32(),
|
|
||||||
error_info: "file not found".to_string(),
|
|
||||||
},
|
|
||||||
rustfs_protos::proto_gen::node_service::Error::default(),
|
|
||||||
],
|
|
||||||
},
|
|
||||||
2,
|
|
||||||
);
|
|
||||||
|
|
||||||
assert!(matches!(errors.as_slice(), [Some(DiskError::FileNotFound), None]));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn delete_versions_response_accepts_legacy_string_errors() {
|
|
||||||
let errors = decode_delete_versions_errors(
|
|
||||||
DeleteVersionsResponse {
|
|
||||||
success: true,
|
|
||||||
errors: vec!["legacy error".to_string(), String::new()],
|
|
||||||
error: None,
|
|
||||||
item_errors: Vec::new(),
|
|
||||||
},
|
|
||||||
2,
|
|
||||||
);
|
|
||||||
|
|
||||||
assert_eq!(errors.len(), 2);
|
|
||||||
assert_eq!(errors[0].as_ref().map(ToString::to_string).as_deref(), Some("io error legacy error"));
|
|
||||||
assert!(errors[1].is_none());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn delete_versions_response_rejects_misaligned_item_errors() {
|
|
||||||
let errors = decode_delete_versions_errors(
|
|
||||||
DeleteVersionsResponse {
|
|
||||||
success: true,
|
|
||||||
errors: vec!["file not found".to_string()],
|
|
||||||
error: None,
|
|
||||||
item_errors: vec![rustfs_protos::proto_gen::node_service::Error {
|
|
||||||
code: DiskError::FileNotFound.to_u32(),
|
|
||||||
error_info: "file not found".to_string(),
|
|
||||||
}],
|
|
||||||
},
|
|
||||||
2,
|
|
||||||
);
|
|
||||||
|
|
||||||
assert_eq!(errors.len(), 2);
|
|
||||||
assert!(errors.iter().all(Option::is_some));
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn disk_mutation_digest_marks_rolling_compatibility() {
|
fn disk_mutation_digest_marks_rolling_compatibility() {
|
||||||
let mut request = Request::new(());
|
let mut request = Request::new(());
|
||||||
|
|||||||
@@ -722,10 +722,6 @@ pub struct DeleteVersionsResponse {
|
|||||||
pub errors: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
|
pub errors: ::prost::alloc::vec::Vec<::prost::alloc::string::String>,
|
||||||
#[prost(message, optional, tag = "3")]
|
#[prost(message, optional, tag = "3")]
|
||||||
pub error: ::core::option::Option<Error>,
|
pub error: ::core::option::Option<Error>,
|
||||||
/// Senders dual-write the legacy strings and typed entries. Receivers prefer typed entries
|
|
||||||
/// when present and fall back to strings for peers that predate this field. Code zero means success.
|
|
||||||
#[prost(message, repeated, tag = "4")]
|
|
||||||
pub item_errors: ::prost::alloc::vec::Vec<Error>,
|
|
||||||
}
|
}
|
||||||
#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
|
#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
|
||||||
pub struct ReadMultipleRequest {
|
pub struct ReadMultipleRequest {
|
||||||
|
|||||||
@@ -493,9 +493,6 @@ message DeleteVersionsResponse {
|
|||||||
bool success = 1;
|
bool success = 1;
|
||||||
repeated string errors = 2;
|
repeated string errors = 2;
|
||||||
optional Error error = 3;
|
optional Error error = 3;
|
||||||
// Senders dual-write the legacy strings and typed entries. Receivers prefer typed entries
|
|
||||||
// when present and fall back to strings for peers that predate this field. Code zero means success.
|
|
||||||
repeated Error item_errors = 4;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
message ReadMultipleRequest {
|
message ReadMultipleRequest {
|
||||||
|
|||||||
@@ -16,7 +16,7 @@ use super::persistence::DataUsageCacheLoadAttempt;
|
|||||||
use super::*;
|
use super::*;
|
||||||
use crate::storage_api::scanner_io::{HTTPRangeSpec, ObjectIO};
|
use crate::storage_api::scanner_io::{HTTPRangeSpec, ObjectIO};
|
||||||
use crate::{ScannerGetObjectReader, ScannerPutObjReader};
|
use crate::{ScannerGetObjectReader, ScannerPutObjReader};
|
||||||
use rustfs_data_usage::{ReplicationAllStats, ReplicationTargetUsage};
|
use rustfs_data_usage::{ReplicationAllStats, ReplicationStats};
|
||||||
use serde_json::Value;
|
use serde_json::Value;
|
||||||
use std::io::Cursor;
|
use std::io::Cursor;
|
||||||
use std::pin::Pin;
|
use std::pin::Pin;
|
||||||
@@ -1636,7 +1636,7 @@ fn size_recursive_prunes_empty_and_preserves_threshold_replication_stats() {
|
|||||||
replication_stats: Some(ReplicationAllStats {
|
replication_stats: Some(ReplicationAllStats {
|
||||||
targets: HashMap::from([(
|
targets: HashMap::from([(
|
||||||
"arn:test:threshold".to_string(),
|
"arn:test:threshold".to_string(),
|
||||||
ReplicationTargetUsage {
|
ReplicationStats {
|
||||||
after_threshold_count: 1,
|
after_threshold_count: 1,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
},
|
},
|
||||||
|
|||||||
@@ -13,7 +13,7 @@
|
|||||||
// limitations under the License.
|
// limitations under the License.
|
||||||
|
|
||||||
use super::*;
|
use super::*;
|
||||||
use rustfs_data_usage::{ReplicationAllStats, ReplicationTargetUsage};
|
use rustfs_data_usage::{ReplicationAllStats, ReplicationStats};
|
||||||
|
|
||||||
const TEST_PLAN_DIGEST: DataUsageScanPlanDigest = DataUsageScanPlanDigest([7; 32]);
|
const TEST_PLAN_DIGEST: DataUsageScanPlanDigest = DataUsageScanPlanDigest([7; 32]);
|
||||||
|
|
||||||
@@ -271,7 +271,7 @@ fn completed_data_usage_info_flattens_nested_bucket_entries() {
|
|||||||
replication_stats: Some(ReplicationAllStats {
|
replication_stats: Some(ReplicationAllStats {
|
||||||
targets: HashMap::from([(
|
targets: HashMap::from([(
|
||||||
"arn:target".to_string(),
|
"arn:target".to_string(),
|
||||||
ReplicationTargetUsage {
|
ReplicationStats {
|
||||||
replicated_size: 2048,
|
replicated_size: 2048,
|
||||||
replicated_count: 2,
|
replicated_count: 2,
|
||||||
..Default::default()
|
..Default::default()
|
||||||
|
|||||||
@@ -146,29 +146,6 @@ fn encode_file_info_msgpack(value: &FileInfo) -> std::result::Result<Vec<u8>, Di
|
|||||||
encode_msgpack_with_capacity(value, "FileInfo", FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT)
|
encode_msgpack_with_capacity(value, "FileInfo", FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn encode_delete_versions_errors(disk_errors: Vec<Option<DiskError>>) -> (Vec<String>, Vec<Error>) {
|
|
||||||
let mut errors = Vec::with_capacity(disk_errors.len());
|
|
||||||
let mut item_errors = Vec::with_capacity(disk_errors.len());
|
|
||||||
for error in disk_errors {
|
|
||||||
match error {
|
|
||||||
Some(error) => {
|
|
||||||
let code = match &error {
|
|
||||||
DiskError::Io(source) if source.kind() == std::io::ErrorKind::NotFound => DiskError::FileNotFound.to_u32(),
|
|
||||||
_ => error.to_u32(),
|
|
||||||
};
|
|
||||||
let error_info = error.to_string();
|
|
||||||
errors.push(error_info.clone());
|
|
||||||
item_errors.push(Error { code, error_info });
|
|
||||||
}
|
|
||||||
None => {
|
|
||||||
errors.push(String::new());
|
|
||||||
item_errors.push(Error::default());
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
(errors, item_errors)
|
|
||||||
}
|
|
||||||
|
|
||||||
fn encode_msgpack_named<T: serde::Serialize>(value: &T, value_name: &str) -> std::result::Result<Vec<u8>, DiskError> {
|
fn encode_msgpack_named<T: serde::Serialize>(value: &T, value_name: &str) -> std::result::Result<Vec<u8>, DiskError> {
|
||||||
let mut serializer = rmp_serde::Serializer::new(Vec::with_capacity(MSGPACK_ENCODE_CAPACITY_HINT)).with_struct_map();
|
let mut serializer = rmp_serde::Serializer::new(Vec::with_capacity(MSGPACK_ENCODE_CAPACITY_HINT)).with_struct_map();
|
||||||
value
|
value
|
||||||
@@ -575,7 +552,6 @@ impl NodeService {
|
|||||||
success: false,
|
success: false,
|
||||||
errors: Vec::new(),
|
errors: Vec::new(),
|
||||||
error: Some(DiskError::other(format!("decode FileInfoVersions failed: {err}")).into()),
|
error: Some(DiskError::other(format!("decode FileInfoVersions failed: {err}")).into()),
|
||||||
item_errors: Vec::new(),
|
|
||||||
}));
|
}));
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
@@ -587,26 +563,30 @@ impl NodeService {
|
|||||||
success: false,
|
success: false,
|
||||||
errors: Vec::new(),
|
errors: Vec::new(),
|
||||||
error: Some(DiskError::other(format!("decode DeleteOptions failed: {err}")).into()),
|
error: Some(DiskError::other(format!("decode DeleteOptions failed: {err}")).into()),
|
||||||
item_errors: Vec::new(),
|
|
||||||
}));
|
}));
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
let (errors, item_errors) =
|
let errors = disk
|
||||||
encode_delete_versions_errors(disk.delete_versions(&request.volume, versions, opts).await);
|
.delete_versions(&request.volume, versions, opts)
|
||||||
|
.await
|
||||||
|
.into_iter()
|
||||||
|
.map(|error| match error {
|
||||||
|
Some(e) => e.to_string(),
|
||||||
|
None => "".to_string(),
|
||||||
|
})
|
||||||
|
.collect();
|
||||||
|
|
||||||
Ok(Response::new(DeleteVersionsResponse {
|
Ok(Response::new(DeleteVersionsResponse {
|
||||||
success: true,
|
success: true,
|
||||||
errors,
|
errors,
|
||||||
error: None,
|
error: None,
|
||||||
item_errors,
|
|
||||||
}))
|
}))
|
||||||
} else {
|
} else {
|
||||||
Ok(Response::new(DeleteVersionsResponse {
|
Ok(Response::new(DeleteVersionsResponse {
|
||||||
success: false,
|
success: false,
|
||||||
errors: Vec::new(),
|
errors: Vec::new(),
|
||||||
error: Some(DiskError::other("cannot find disk".to_string()).into()),
|
error: Some(DiskError::other("cannot find disk".to_string()).into()),
|
||||||
item_errors: Vec::new(),
|
|
||||||
}))
|
}))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -1632,8 +1612,8 @@ impl NodeService {
|
|||||||
mod tests {
|
mod tests {
|
||||||
use super::{
|
use super::{
|
||||||
compat_response_json, decode_msgpack_or_json, decode_rename_data_request_file_info,
|
compat_response_json, decode_msgpack_or_json, decode_rename_data_request_file_info,
|
||||||
encode_batch_read_version_response_payloads, encode_delete_versions_errors, encode_file_info_msgpack, encode_msgpack,
|
encode_batch_read_version_response_payloads, encode_file_info_msgpack, encode_msgpack, encode_msgpack_named,
|
||||||
encode_msgpack_named, encode_read_multiple_response_payloads, encode_rename_data_response_payloads,
|
encode_read_multiple_response_payloads, encode_rename_data_response_payloads,
|
||||||
};
|
};
|
||||||
use crate::storage::rpc::node_service::make_server;
|
use crate::storage::rpc::node_service::make_server;
|
||||||
use crate::storage::storage_api::ReadMultipleResp;
|
use crate::storage::storage_api::ReadMultipleResp;
|
||||||
@@ -1652,18 +1632,6 @@ mod tests {
|
|||||||
count: u32,
|
count: u32,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn delete_versions_response_dual_writes_typed_item_errors() {
|
|
||||||
let raw_not_found = super::DiskError::Io(std::io::Error::from(std::io::ErrorKind::NotFound));
|
|
||||||
let (errors, item_errors) = encode_delete_versions_errors(vec![Some(raw_not_found), None]);
|
|
||||||
|
|
||||||
assert!(errors[0].starts_with("io error "));
|
|
||||||
assert!(errors[1].is_empty());
|
|
||||||
assert_eq!(item_errors[0].code, super::DiskError::FileNotFound.to_u32());
|
|
||||||
assert_eq!(item_errors[0].error_info, errors[0]);
|
|
||||||
assert_eq!(item_errors[1].code, 0);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[serial]
|
#[serial]
|
||||||
async fn handle_read_version_records_attribution_for_missing_disk() {
|
async fn handle_read_version_records_attribution_for_missing_disk() {
|
||||||
|
|||||||
Reference in New Issue
Block a user