feat(ecstore): coalesce GET ReadVersion RPCs (#6395)

This commit is contained in:
houseme
2026-08-23 01:15:38 +08:00
committed by GitHub
parent 7c1a76dfd9
commit f1b92af4a3
13 changed files with 987 additions and 74 deletions
+5
View File
@@ -20,6 +20,7 @@ use crate::storage_api::server::readiness::contract::admin::StorageAdminApi;
use crate::storage_api::server::readiness::{Endpoint, EndpointServerPools, is_dist_erasure};
#[cfg(test)]
use crate::storage_api::server::readiness::{Endpoints, PoolEndpoints};
use crate::storage_api::startup::shutdown::mark_get_metadata_read_version_coalescing_service_ready;
use bytes::Bytes;
use http::HeaderValue;
use http::{Request as HttpRequest, Response, StatusCode};
@@ -212,6 +213,9 @@ where
if readiness_gate_blocks_path(path, &readiness) {
return Ok(service_not_ready_response(readiness.current_stage()));
}
if !is_probe_path(path) && readiness.is_ready() {
mark_get_metadata_read_version_coalescing_service_ready();
}
let resp = inner.call(req).await?;
// System is ready, forward to the actual S3/RPC handlers
// Transparently converts any response body into a BoxBody, and then Trace/Cors/Compression continues to work
@@ -232,6 +236,7 @@ pub async fn publish_ready_when_runtime_ready(
collect_node_readiness,
|dependency_readiness| {
readiness.mark_stage(rustfs_common::SystemStage::FullReady);
mark_get_metadata_read_version_coalescing_service_ready();
if let Some(state_manager) = state_manager {
state_manager.update(ServiceState::Ready);
}
+99 -15
View File
@@ -23,10 +23,12 @@ use bytes::Bytes;
use rustfs_filemeta::FileInfo;
use rustfs_io_metrics::internode_metrics::{
INTERNODE_MSGPACK_CODEC_JSON, INTERNODE_MSGPACK_CODEC_MSGPACK, INTERNODE_MSGPACK_DIRECTION_REQUEST,
INTERNODE_OPERATION_GRPC_READ_ALL, INTERNODE_OPERATION_GRPC_READ_VERSION, INTERNODE_OPERATION_GRPC_WRITE_ALL,
INTERNODE_STAGE_READ_VERSION_DISK_READ, INTERNODE_STAGE_READ_VERSION_REQUEST_DECODE,
INTERNODE_STAGE_READ_VERSION_RESPONSE_JSON_ENCODE, INTERNODE_STAGE_READ_VERSION_RESPONSE_MSGPACK_ENCODE,
INTERNODE_TRANSPORT_BACKEND_GRPC, global_internode_metrics,
INTERNODE_OPERATION_GRPC_BATCH_READ_VERSION, INTERNODE_OPERATION_GRPC_READ_ALL, INTERNODE_OPERATION_GRPC_READ_VERSION,
INTERNODE_OPERATION_GRPC_WRITE_ALL, INTERNODE_STAGE_BATCH_READ_VERSION_DISK_READ,
INTERNODE_STAGE_BATCH_READ_VERSION_REQUEST_DECODE, INTERNODE_STAGE_BATCH_READ_VERSION_RESPONSE_JSON_ENCODE,
INTERNODE_STAGE_BATCH_READ_VERSION_RESPONSE_MSGPACK_ENCODE, INTERNODE_STAGE_READ_VERSION_DISK_READ,
INTERNODE_STAGE_READ_VERSION_REQUEST_DECODE, INTERNODE_STAGE_READ_VERSION_RESPONSE_JSON_ENCODE,
INTERNODE_STAGE_READ_VERSION_RESPONSE_MSGPACK_ENCODE, INTERNODE_TRANSPORT_BACKEND_GRPC, global_internode_metrics,
};
use rustfs_protos::proto_gen::node_service::*;
use serde::de::DeserializeOwned;
@@ -242,24 +244,42 @@ fn record_read_version_stage(stage: &'static str, started_at: Option<Instant>) {
}
}
fn record_batch_read_version_stage(stage: &'static str, started_at: Option<Instant>) {
if let Some(started_at) = started_at {
global_internode_metrics().record_stage_duration_for_operation_and_backend(
INTERNODE_OPERATION_GRPC_BATCH_READ_VERSION,
INTERNODE_TRANSPORT_BACKEND_GRPC,
stage,
started_at.elapsed(),
);
}
}
fn encode_batch_read_version_response_payloads(
batch_read_version_resps: &[BatchReadVersionResp],
request_decoded_from_msgpack: bool,
) -> std::result::Result<(Vec<String>, Vec<Bytes>), DiskError> {
let attribution_enabled = rustfs_io_metrics::get_stage_metrics_enabled();
let mut batch_read_version_resps_json = Vec::with_capacity(batch_read_version_resps.len());
let mut batch_read_version_resps_bin = Vec::with_capacity(batch_read_version_resps.len());
let json_encode_started = internode_stage_timer(attribution_enabled);
for batch_read_version_resp in batch_read_version_resps {
batch_read_version_resps_json.push(
compat_response_json(batch_read_version_resp, request_decoded_from_msgpack)
.map_err(|err| DiskError::other(format!("encode BatchReadVersionResp json failed: {err}")))?,
);
}
record_batch_read_version_stage(INTERNODE_STAGE_BATCH_READ_VERSION_RESPONSE_JSON_ENCODE, json_encode_started);
let mut batch_read_version_resps_bin = Vec::with_capacity(batch_read_version_resps.len());
let msgpack_encode_started = internode_stage_timer(attribution_enabled);
for batch_read_version_resp in batch_read_version_resps {
batch_read_version_resps_bin.push(Bytes::from(encode_msgpack_with_capacity(
batch_read_version_resp,
"BatchReadVersionResp",
FILE_INFO_MSGPACK_ENCODE_CAPACITY_HINT,
)?));
}
record_batch_read_version_stage(INTERNODE_STAGE_BATCH_READ_VERSION_RESPONSE_MSGPACK_ENCODE, msgpack_encode_started);
Ok((batch_read_version_resps_json, batch_read_version_resps_bin))
}
@@ -485,15 +505,37 @@ impl NodeService {
&self,
request: Request<BatchReadVersionRequest>,
) -> Result<Response<BatchReadVersionResponse>, Status> {
let attribution_enabled = rustfs_io_metrics::get_stage_metrics_enabled();
let request = request.into_inner();
if attribution_enabled {
let metrics = global_internode_metrics();
metrics.record_incoming_request_for_operation_and_backend(
INTERNODE_OPERATION_GRPC_BATCH_READ_VERSION,
INTERNODE_TRANSPORT_BACKEND_GRPC,
);
metrics.record_recv_bytes_for_operation_and_backend(
INTERNODE_OPERATION_GRPC_BATCH_READ_VERSION,
INTERNODE_TRANSPORT_BACKEND_GRPC,
request
.disk
.len()
.saturating_add(request.batch_read_version_req.len())
.saturating_add(request.batch_read_version_req_bin.len()),
);
}
if let Some(disk) = self.find_disk(&request.disk).await {
let decode_started = internode_stage_timer(attribution_enabled);
let decoded_batch_read_version_req: DecodedRpcPayload<BatchReadVersionReq> = match decode_msgpack_or_json_with_source(
&request.batch_read_version_req_bin,
&request.batch_read_version_req,
"BatchReadVersionReq",
) {
Ok(batch_read_version_req) => batch_read_version_req,
Ok(batch_read_version_req) => {
record_batch_read_version_stage(INTERNODE_STAGE_BATCH_READ_VERSION_REQUEST_DECODE, decode_started);
batch_read_version_req
}
Err(err) => {
record_batch_read_version_stage(INTERNODE_STAGE_BATCH_READ_VERSION_REQUEST_DECODE, decode_started);
return Ok(Response::new(BatchReadVersionResponse {
success: false,
batch_read_version_resps: Vec::new(),
@@ -514,8 +556,10 @@ impl NodeService {
}));
}
let disk_read_started = internode_stage_timer(attribution_enabled);
match disk.batch_read_version(batch_read_version_req).await {
Ok(batch_read_version_resps) => {
record_batch_read_version_stage(INTERNODE_STAGE_BATCH_READ_VERSION_DISK_READ, disk_read_started);
let (batch_read_version_resps, batch_read_version_resps_bin) =
match encode_batch_read_version_response_payloads(&batch_read_version_resps, request_decoded_from_msgpack)
{
@@ -537,12 +581,15 @@ impl NodeService {
error: None,
}))
}
Err(err) => Ok(Response::new(BatchReadVersionResponse {
success: false,
batch_read_version_resps: Vec::new(),
batch_read_version_resps_bin: Vec::new(),
error: Some(err.into()),
})),
Err(err) => {
record_batch_read_version_stage(INTERNODE_STAGE_BATCH_READ_VERSION_DISK_READ, disk_read_started);
Ok(Response::new(BatchReadVersionResponse {
success: false,
batch_read_version_resps: Vec::new(),
batch_read_version_resps_bin: Vec::new(),
error: Some(err.into()),
}))
}
}
} else {
Ok(Response::new(BatchReadVersionResponse {
@@ -722,9 +769,9 @@ impl NodeService {
&self,
request: Request<ReadVersionRequest>,
) -> Result<Response<ReadVersionResponse>, Status> {
let request = request.into_inner();
let metrics = global_internode_metrics();
let read_version_attribution_enabled = rustfs_io_metrics::get_stage_metrics_enabled();
let request = request.into_inner();
if read_version_attribution_enabled {
metrics.record_incoming_request_for_operation_and_backend(
INTERNODE_OPERATION_GRPC_READ_VERSION,
@@ -1635,6 +1682,7 @@ mod tests {
encode_batch_read_version_response_payloads, encode_delete_versions_errors, encode_file_info_msgpack, encode_msgpack,
encode_msgpack_named, encode_read_multiple_response_payloads, encode_rename_data_response_payloads,
};
use crate::storage::DiskError;
use crate::storage::rpc::node_service::make_server;
use crate::storage::storage_api::ReadMultipleResp;
use crate::storage::storage_api::RenameDataResp;
@@ -2028,7 +2076,8 @@ mod tests {
path: "object-a".to_string(),
version_id: "version-a".to_string(),
success: false,
error: "file version not found".to_string(),
error: DiskError::FileVersionNotFound.to_string(),
error_code: DiskError::FileVersionNotFound.to_u32(),
..Default::default()
}];
@@ -2044,8 +2093,43 @@ mod tests {
.expect("msgpack batch read version response should decode");
assert_eq!(json_decoded.index, responses[0].index);
assert_eq!(json_decoded.error_code, responses[0].error_code);
assert_eq!(msgpack_decoded.path, responses[0].path);
assert_eq!(msgpack_decoded.error, responses[0].error);
assert_eq!(msgpack_decoded.error_code, responses[0].error_code);
}
#[test]
fn batch_read_version_response_decode_accepts_legacy_payload_without_error_code() {
#[derive(Serialize)]
struct LegacyBatchReadVersionResp {
index: usize,
path: String,
version_id: String,
success: bool,
file_info: FileInfo,
error: String,
}
let legacy = LegacyBatchReadVersionResp {
index: 2,
path: "object-legacy".to_string(),
version_id: "version-legacy".to_string(),
success: false,
file_info: FileInfo::default(),
error: "legacy error".to_string(),
};
let legacy_json = serde_json::to_string(&legacy).expect("legacy json should encode");
let legacy_msgpack = encode_msgpack(&legacy, "LegacyBatchReadVersionResp").expect("legacy msgpack should encode");
let json_decoded: BatchReadVersionResp =
decode_msgpack_or_json(&[], &legacy_json, "BatchReadVersionResp").expect("legacy json should decode");
let msgpack_decoded: BatchReadVersionResp =
decode_msgpack_or_json(&legacy_msgpack, "", "BatchReadVersionResp").expect("legacy msgpack should decode");
assert_eq!(json_decoded.error_code, 0);
assert_eq!(msgpack_decoded.error_code, 0);
assert_eq!(msgpack_decoded.error, legacy.error);
}
#[test]
+4
View File
@@ -1100,6 +1100,10 @@ pub(crate) fn shutdown_background_monitors() {
rustfs_ecstore::shutdown_background_monitors();
}
pub(crate) fn mark_get_metadata_read_version_coalescing_service_ready() {
rustfs_ecstore::mark_get_metadata_read_version_coalescing_service_ready();
}
pub(crate) fn set_global_rustfs_port(value: u16) {
ecstore_global::set_global_rustfs_port(value);
}
+2 -1
View File
@@ -284,7 +284,8 @@ pub(crate) mod startup {
pub(crate) mod shutdown {
pub(crate) use crate::storage::storage_api::{
shutdown_background_monitors, shutdown_background_services, store_compression_total_in_backend,
mark_get_metadata_read_version_coalescing_service_ready, shutdown_background_monitors, shutdown_background_services,
store_compression_total_in_backend,
};
}