mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-29 08:27:06 +00:00
perf(storage): skip read-version JSON for bin peers (#5808)
Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -2442,7 +2442,7 @@ mod tests {
|
|||||||
json_field: "batch_read_version_resps",
|
json_field: "batch_read_version_resps",
|
||||||
bin_field: "batch_read_version_resps_bin",
|
bin_field: "batch_read_version_resps_bin",
|
||||||
},
|
},
|
||||||
json_encoder: "compat_response_json(batch_read_version_resp)",
|
json_encoder: "compat_response_json(batch_read_version_resp, false)",
|
||||||
},
|
},
|
||||||
ResponseCompatSendSite {
|
ResponseCompatSendSite {
|
||||||
field: CompatPayloadField {
|
field: CompatPayloadField {
|
||||||
@@ -2450,7 +2450,7 @@ mod tests {
|
|||||||
json_field: "read_multiple_resps",
|
json_field: "read_multiple_resps",
|
||||||
bin_field: "read_multiple_resps_bin",
|
bin_field: "read_multiple_resps_bin",
|
||||||
},
|
},
|
||||||
json_encoder: "compat_response_json(read_multiple_resp)",
|
json_encoder: "compat_response_json(read_multiple_resp, false)",
|
||||||
},
|
},
|
||||||
ResponseCompatSendSite {
|
ResponseCompatSendSite {
|
||||||
field: CompatPayloadField {
|
field: CompatPayloadField {
|
||||||
@@ -2458,7 +2458,7 @@ mod tests {
|
|||||||
json_field: "file_info",
|
json_field: "file_info",
|
||||||
bin_field: "file_info_bin",
|
bin_field: "file_info_bin",
|
||||||
},
|
},
|
||||||
json_encoder: "let file_info_json = compat_response_json(&file_info);",
|
json_encoder: "let file_info_json = compat_response_json(&file_info, request_had_msgpack_payload);",
|
||||||
},
|
},
|
||||||
ResponseCompatSendSite {
|
ResponseCompatSendSite {
|
||||||
field: CompatPayloadField {
|
field: CompatPayloadField {
|
||||||
@@ -2466,7 +2466,7 @@ mod tests {
|
|||||||
json_field: "raw_file_info",
|
json_field: "raw_file_info",
|
||||||
bin_field: "raw_file_info_bin",
|
bin_field: "raw_file_info_bin",
|
||||||
},
|
},
|
||||||
json_encoder: "let raw_file_info_json = compat_response_json(&raw_file_info);",
|
json_encoder: "let raw_file_info_json = compat_response_json(&raw_file_info, false);",
|
||||||
},
|
},
|
||||||
ResponseCompatSendSite {
|
ResponseCompatSendSite {
|
||||||
field: CompatPayloadField {
|
field: CompatPayloadField {
|
||||||
@@ -2474,7 +2474,7 @@ mod tests {
|
|||||||
json_field: "rename_data_resp",
|
json_field: "rename_data_resp",
|
||||||
bin_field: "rename_data_resp_bin",
|
bin_field: "rename_data_resp_bin",
|
||||||
},
|
},
|
||||||
json_encoder: "let rename_data_resp_json = compat_response_json(&rename_data_resp);",
|
json_encoder: "let rename_data_resp_json = compat_response_json(&rename_data_resp, false);",
|
||||||
},
|
},
|
||||||
];
|
];
|
||||||
|
|
||||||
|
|||||||
@@ -141,11 +141,16 @@ fn verify_disk_mutation_digest<T>(
|
|||||||
.map_err(|err| Status::permission_denied(format!("{op} authentication failed: {err}")))
|
.map_err(|err| Status::permission_denied(format!("{op} authentication failed: {err}")))
|
||||||
}
|
}
|
||||||
|
|
||||||
/// JSON compatibility string for a dual-encoded response field. Returns an empty string only when
|
/// JSON compatibility string for a dual-encoded response field. Returns an empty string when
|
||||||
/// msgpack-only mode and its explicit fleet confirmation guard are both enabled; otherwise the
|
/// msgpack-only mode and its explicit fleet confirmation guard are both enabled, or when
|
||||||
/// legacy JSON encoding is retained for old peers. The paired `_bin` field is always sent.
|
/// request-local `_bin` payloads prove the caller can consume the paired response `_bin` field,
|
||||||
fn compat_response_json<T: serde::Serialize>(value: &T) -> std::result::Result<String, serde_json::Error> {
|
/// while legacy JSON-only callers still need the compatibility response string during rollout.
|
||||||
if rustfs_protos::internode_rpc_msgpack_only() {
|
/// The paired `_bin` field is always sent.
|
||||||
|
fn compat_response_json<T: serde::Serialize>(
|
||||||
|
value: &T,
|
||||||
|
request_had_msgpack_payload: bool,
|
||||||
|
) -> std::result::Result<String, serde_json::Error> {
|
||||||
|
if request_had_msgpack_payload || rustfs_protos::internode_rpc_msgpack_only() {
|
||||||
return Ok(String::new());
|
return Ok(String::new());
|
||||||
}
|
}
|
||||||
serde_json::to_string(value)
|
serde_json::to_string(value)
|
||||||
@@ -159,7 +164,7 @@ fn encode_read_multiple_response_payloads(
|
|||||||
|
|
||||||
for read_multiple_resp in read_multiple_resps {
|
for read_multiple_resp in read_multiple_resps {
|
||||||
read_multiple_resps_json.push(
|
read_multiple_resps_json.push(
|
||||||
compat_response_json(read_multiple_resp)
|
compat_response_json(read_multiple_resp, false)
|
||||||
.map_err(|err| DiskError::other(format!("encode ReadMultipleResp json failed: {err}")))?,
|
.map_err(|err| DiskError::other(format!("encode ReadMultipleResp json failed: {err}")))?,
|
||||||
);
|
);
|
||||||
read_multiple_resps_bin.push(Bytes::from(encode_msgpack(read_multiple_resp, "ReadMultipleResp")?));
|
read_multiple_resps_bin.push(Bytes::from(encode_msgpack(read_multiple_resp, "ReadMultipleResp")?));
|
||||||
@@ -176,7 +181,7 @@ fn encode_batch_read_version_response_payloads(
|
|||||||
|
|
||||||
for batch_read_version_resp in batch_read_version_resps {
|
for batch_read_version_resp in batch_read_version_resps {
|
||||||
batch_read_version_resps_json.push(
|
batch_read_version_resps_json.push(
|
||||||
compat_response_json(batch_read_version_resp)
|
compat_response_json(batch_read_version_resp, false)
|
||||||
.map_err(|err| DiskError::other(format!("encode BatchReadVersionResp json failed: {err}")))?,
|
.map_err(|err| DiskError::other(format!("encode BatchReadVersionResp json failed: {err}")))?,
|
||||||
);
|
);
|
||||||
batch_read_version_resps_bin.push(Bytes::from(encode_msgpack_with_capacity(
|
batch_read_version_resps_bin.push(Bytes::from(encode_msgpack_with_capacity(
|
||||||
@@ -571,7 +576,7 @@ impl NodeService {
|
|||||||
if let Some(disk) = self.find_disk(&request.disk).await {
|
if let Some(disk) = self.find_disk(&request.disk).await {
|
||||||
match disk.read_xl(&request.volume, &request.path, request.read_data).await {
|
match disk.read_xl(&request.volume, &request.path, request.read_data).await {
|
||||||
Ok(raw_file_info) => {
|
Ok(raw_file_info) => {
|
||||||
let raw_file_info_json = compat_response_json(&raw_file_info);
|
let raw_file_info_json = compat_response_json(&raw_file_info, false);
|
||||||
let raw_file_info_bin = encode_msgpack(&raw_file_info, "RawFileInfo");
|
let raw_file_info_bin = encode_msgpack(&raw_file_info, "RawFileInfo");
|
||||||
match (raw_file_info_json, raw_file_info_bin) {
|
match (raw_file_info_json, raw_file_info_bin) {
|
||||||
(Ok(raw_file_info), Ok(raw_file_info_bin)) => Ok(Response::new(ReadXlResponse {
|
(Ok(raw_file_info), Ok(raw_file_info_bin)) => Ok(Response::new(ReadXlResponse {
|
||||||
@@ -617,6 +622,7 @@ impl NodeService {
|
|||||||
) -> Result<Response<ReadVersionResponse>, Status> {
|
) -> Result<Response<ReadVersionResponse>, Status> {
|
||||||
let request = request.into_inner();
|
let request = request.into_inner();
|
||||||
if let Some(disk) = self.find_disk(&request.disk).await {
|
if let Some(disk) = self.find_disk(&request.disk).await {
|
||||||
|
let request_had_msgpack_payload = !request.opts_bin.is_empty();
|
||||||
let opts = match decode_msgpack_or_json::<ReadOptions>(&request.opts_bin, &request.opts, "ReadOptions") {
|
let opts = match decode_msgpack_or_json::<ReadOptions>(&request.opts_bin, &request.opts, "ReadOptions") {
|
||||||
Ok(options) => options,
|
Ok(options) => options,
|
||||||
Err(err) => {
|
Err(err) => {
|
||||||
@@ -633,7 +639,7 @@ impl NodeService {
|
|||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
Ok(file_info) => {
|
Ok(file_info) => {
|
||||||
let file_info_json = compat_response_json(&file_info);
|
let file_info_json = compat_response_json(&file_info, request_had_msgpack_payload);
|
||||||
let file_info_bin = encode_file_info_msgpack(&file_info);
|
let file_info_bin = encode_file_info_msgpack(&file_info);
|
||||||
match (file_info_json, file_info_bin) {
|
match (file_info_json, file_info_bin) {
|
||||||
(Ok(file_info), Ok(file_info_bin)) => Ok(Response::new(ReadVersionResponse {
|
(Ok(file_info), Ok(file_info_bin)) => Ok(Response::new(ReadVersionResponse {
|
||||||
@@ -979,7 +985,7 @@ impl NodeService {
|
|||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
Ok(rename_data_resp) => {
|
Ok(rename_data_resp) => {
|
||||||
let rename_data_resp_json = compat_response_json(&rename_data_resp);
|
let rename_data_resp_json = compat_response_json(&rename_data_resp, false);
|
||||||
let rename_data_resp_bin = encode_msgpack_named(&rename_data_resp, "RenameDataResp");
|
let rename_data_resp_bin = encode_msgpack_named(&rename_data_resp, "RenameDataResp");
|
||||||
match (rename_data_resp_json, rename_data_resp_bin) {
|
match (rename_data_resp_json, rename_data_resp_bin) {
|
||||||
(Ok(rename_data_resp), Ok(rename_data_resp_bin)) => Ok(Response::new(RenameDataResponse {
|
(Ok(rename_data_resp), Ok(rename_data_resp_bin)) => Ok(Response::new(RenameDataResponse {
|
||||||
@@ -1527,7 +1533,7 @@ mod tests {
|
|||||||
name: "legacy-json".to_string(),
|
name: "legacy-json".to_string(),
|
||||||
count: 9,
|
count: 9,
|
||||||
};
|
};
|
||||||
let json = compat_response_json(&payload).expect("compat response json should encode");
|
let json = compat_response_json(&payload, false).expect("compat response json should encode");
|
||||||
|
|
||||||
assert!(!json.is_empty(), "old JSON clients must remain compatible without fleet confirmation");
|
assert!(!json.is_empty(), "old JSON clients must remain compatible without fleet confirmation");
|
||||||
assert_eq!(json, serde_json::to_string(&payload).expect("json should encode"));
|
assert_eq!(json, serde_json::to_string(&payload).expect("json should encode"));
|
||||||
@@ -1547,7 +1553,7 @@ mod tests {
|
|||||||
name: "msgpack-only".to_string(),
|
name: "msgpack-only".to_string(),
|
||||||
count: 10,
|
count: 10,
|
||||||
};
|
};
|
||||||
let json = compat_response_json(&payload).expect("compat response json should encode");
|
let json = compat_response_json(&payload, false).expect("compat response json should encode");
|
||||||
|
|
||||||
assert!(json.is_empty(), "msgpack-only may empty JSON only after explicit fleet confirmation");
|
assert!(json.is_empty(), "msgpack-only may empty JSON only after explicit fleet confirmation");
|
||||||
},
|
},
|
||||||
@@ -1567,7 +1573,7 @@ mod tests {
|
|||||||
(rustfs_config::ENV_INTERNODE_RPC_MSGPACK_ONLY_FLEET_CONFIRMED, Some("true")),
|
(rustfs_config::ENV_INTERNODE_RPC_MSGPACK_ONLY_FLEET_CONFIRMED, Some("true")),
|
||||||
],
|
],
|
||||||
|| {
|
|| {
|
||||||
let json = compat_response_json(&payload).expect("compat response json should encode");
|
let json = compat_response_json(&payload, false).expect("compat response json should encode");
|
||||||
|
|
||||||
assert!(json.is_empty(), "both gates should enter msgpack-only response mode");
|
assert!(json.is_empty(), "both gates should enter msgpack-only response mode");
|
||||||
},
|
},
|
||||||
@@ -1584,7 +1590,7 @@ mod tests {
|
|||||||
],
|
],
|
||||||
] {
|
] {
|
||||||
with_internode_msgpack_env(vars, || {
|
with_internode_msgpack_env(vars, || {
|
||||||
let json = compat_response_json(&payload).expect("compat response json should encode");
|
let json = compat_response_json(&payload, false).expect("compat response json should encode");
|
||||||
|
|
||||||
assert!(!json.is_empty(), "removing either gate should restore old-peer JSON compatibility");
|
assert!(!json.is_empty(), "removing either gate should restore old-peer JSON compatibility");
|
||||||
assert_eq!(json, serde_json::to_string(&payload).expect("json should encode"));
|
assert_eq!(json, serde_json::to_string(&payload).expect("json should encode"));
|
||||||
@@ -1592,6 +1598,45 @@ mod tests {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn compat_response_json_omits_json_for_msgpack_capable_request_without_fleet_gate() {
|
||||||
|
with_internode_msgpack_env(
|
||||||
|
[
|
||||||
|
(rustfs_config::ENV_INTERNODE_RPC_MSGPACK_ONLY, None::<&str>),
|
||||||
|
(rustfs_config::ENV_INTERNODE_RPC_MSGPACK_ONLY_FLEET_CONFIRMED, None::<&str>),
|
||||||
|
],
|
||||||
|
|| {
|
||||||
|
let payload = SamplePayload {
|
||||||
|
name: "request-local-bin".to_string(),
|
||||||
|
count: 13,
|
||||||
|
};
|
||||||
|
let json = compat_response_json(&payload, true).expect("compat response json should encode");
|
||||||
|
|
||||||
|
assert!(json.is_empty(), "a request-local msgpack payload may receive a bin-only response");
|
||||||
|
},
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn compat_response_json_keeps_json_for_legacy_json_request() {
|
||||||
|
with_internode_msgpack_env(
|
||||||
|
[
|
||||||
|
(rustfs_config::ENV_INTERNODE_RPC_MSGPACK_ONLY, None::<&str>),
|
||||||
|
(rustfs_config::ENV_INTERNODE_RPC_MSGPACK_ONLY_FLEET_CONFIRMED, None::<&str>),
|
||||||
|
],
|
||||||
|
|| {
|
||||||
|
let payload = SamplePayload {
|
||||||
|
name: "legacy-request-json".to_string(),
|
||||||
|
count: 14,
|
||||||
|
};
|
||||||
|
let json = compat_response_json(&payload, false).expect("compat response json should encode");
|
||||||
|
|
||||||
|
assert!(!json.is_empty(), "legacy JSON-only callers must keep receiving response JSON");
|
||||||
|
assert_eq!(json, serde_json::to_string(&payload).expect("json should encode"));
|
||||||
|
},
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn decode_msgpack_or_json_fails_closed_on_corrupt_non_empty_msgpack() {
|
fn decode_msgpack_or_json_fails_closed_on_corrupt_non_empty_msgpack() {
|
||||||
let before = global_internode_metrics().msgpack_json_decode_error_total_for_test();
|
let before = global_internode_metrics().msgpack_json_decode_error_total_for_test();
|
||||||
|
|||||||
Reference in New Issue
Block a user