refactor: centralize network client runtime sources (#3796)

This commit is contained in:
Zhengchao An
2026-06-23 22:35:51 +08:00
committed by GitHub
parent 726c26fa01
commit e59e1852ec
9 changed files with 210 additions and 75 deletions
+1
View File
@@ -35,6 +35,7 @@ pub mod constants;
pub mod credentials;
pub mod object_api_utils;
pub mod object_handlers_common;
pub(crate) mod runtime_sources;
pub mod signer_error;
pub mod transition_api;
pub mod utils;
@@ -0,0 +1,25 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use rustfs_tls_runtime::{GlobalPublishedOutboundTlsState, load_global_outbound_tls_state, record_tls_generation};
const ECSTORE_TRANSITION_CLIENT_TLS_CONSUMER: &str = "ecstore_transition_client";
pub(crate) async fn transition_client_outbound_tls_state() -> GlobalPublishedOutboundTlsState {
load_global_outbound_tls_state().await
}
pub(crate) fn record_transition_client_tls_generation(generation: u64) {
record_tls_generation(ECSTORE_TRANSITION_CLIENT_TLS_CONSUMER, generation);
}
+2 -3
View File
@@ -51,7 +51,6 @@ use md5::Md5;
use rand::{Rng, RngExt};
use rustfs_config::MAX_S3_CLIENT_RESPONSE_SIZE;
use rustfs_rio::HashReader;
use rustfs_tls_runtime::{load_global_outbound_tls_state, record_tls_generation};
use rustfs_utils::HashAlgorithm;
use rustfs_utils::{
net::get_endpoint_url,
@@ -192,8 +191,8 @@ where
async fn build_tls_config() -> Result<rustls::ClientConfig, std::io::Error> {
with_rustls_init_guard(|| Ok(()))?;
let outbound_tls = load_global_outbound_tls_state().await;
record_tls_generation("ecstore_transition_client", outbound_tls.generation.0);
let outbound_tls = crate::client::runtime_sources::transition_client_outbound_tls_state().await;
crate::client::runtime_sources::record_transition_client_tls_generation(outbound_tls.generation.0);
let builder = if let Some(root_ca_pem) = outbound_tls.root_ca_pem.as_ref() {
let mut reader = std::io::BufReader::new(root_ca_pem.as_slice());
let certs_der = rustls_pki_types::CertificateDer::pem_reader_iter(&mut reader)
+1
View File
@@ -20,6 +20,7 @@ mod peer_rest_client;
mod peer_s3_client;
mod remote_disk;
mod remote_locker;
mod runtime_sources;
pub use client::{
TonicInterceptor, gen_tonic_signature_interceptor, node_service_time_out_client, node_service_time_out_client_no_auth,
+14 -58
View File
@@ -34,10 +34,6 @@ use bytes::Bytes;
use futures::lock::Mutex;
use metrics::counter;
use rustfs_filemeta::{FileInfo, ObjectPartInfo, RawFileInfo};
use rustfs_io_metrics::internode_metrics::{
INTERNODE_OPERATION_GRPC_READ_ALL, INTERNODE_OPERATION_GRPC_WRITE_ALL, INTERNODE_TRANSPORT_BACKEND_GRPC,
INTERNODE_TRANSPORT_BACKEND_TCP_HTTP, global_internode_metrics,
};
use rustfs_protos::evict_failed_connection;
use rustfs_protos::proto_gen::node_service::RenamePartRequest;
use rustfs_protos::proto_gen::node_service::{
@@ -187,22 +183,14 @@ impl RemoteDisk {
if attempt > 1
&& let Some(classification) = last_retry_classification
{
global_internode_metrics().record_retry_success_for_operation_and_backend(
rustfs_io_metrics::internode_metrics::INTERNODE_OPERATION_PUT_FILE_STREAM,
INTERNODE_TRANSPORT_BACKEND_TCP_HTTP,
classification,
);
crate::rpc::runtime_sources::record_remote_disk_open_write_retry_success(classification);
}
return Ok(writer);
}
Err(err) if attempt < REMOTE_DISK_OPEN_WRITE_MAX_ATTEMPTS && Self::is_retryable_open_write_error(&err) => {
if let Some(classification) = err.internode_http_error_kind() {
let classification = classification.metric_label();
global_internode_metrics().record_retry_for_operation_and_backend(
rustfs_io_metrics::internode_metrics::INTERNODE_OPERATION_PUT_FILE_STREAM,
INTERNODE_TRANSPORT_BACKEND_TCP_HTTP,
classification,
);
crate::rpc::runtime_sources::record_remote_disk_open_write_retry(classification);
last_retry_classification = Some(classification);
}
debug!(
@@ -2109,10 +2097,7 @@ impl DiskAPI for RemoteDisk {
let data_len = data.len();
let disk = self.disk_ref().await;
let mut client = self.get_client().await.map_err(|err| {
global_internode_metrics().record_error_for_operation_and_backend(
INTERNODE_OPERATION_GRPC_WRITE_ALL,
INTERNODE_TRANSPORT_BACKEND_GRPC,
);
crate::rpc::runtime_sources::record_remote_disk_grpc_write_all_error();
Error::other(format!("can not get client, err: {err}"))
})?;
let request = Request::new(WriteAllRequest {
@@ -2122,32 +2107,19 @@ impl DiskAPI for RemoteDisk {
data,
});
global_internode_metrics().record_outgoing_request_for_operation_and_backend(
INTERNODE_OPERATION_GRPC_WRITE_ALL,
INTERNODE_TRANSPORT_BACKEND_GRPC,
);
crate::rpc::runtime_sources::record_remote_disk_grpc_write_all_request();
let response = match client.write_all(request).await {
Ok(response) => response.into_inner(),
Err(err) => {
global_internode_metrics().record_error_for_operation_and_backend(
INTERNODE_OPERATION_GRPC_WRITE_ALL,
INTERNODE_TRANSPORT_BACKEND_GRPC,
);
crate::rpc::runtime_sources::record_remote_disk_grpc_write_all_error();
return Err(err.into());
}
};
global_internode_metrics().record_sent_bytes_for_operation_and_backend(
INTERNODE_OPERATION_GRPC_WRITE_ALL,
INTERNODE_TRANSPORT_BACKEND_GRPC,
data_len,
);
crate::rpc::runtime_sources::record_remote_disk_grpc_write_all_sent_bytes(data_len);
if !response.success {
global_internode_metrics().record_error_for_operation_and_backend(
INTERNODE_OPERATION_GRPC_WRITE_ALL,
INTERNODE_TRANSPORT_BACKEND_GRPC,
);
crate::rpc::runtime_sources::record_remote_disk_grpc_write_all_error();
return Err(response.error.unwrap_or_default().into());
}
@@ -2176,10 +2148,7 @@ impl DiskAPI for RemoteDisk {
|| async {
let disk = self.disk_ref().await;
let mut client = self.get_client().await.map_err(|err| {
global_internode_metrics().record_error_for_operation_and_backend(
INTERNODE_OPERATION_GRPC_READ_ALL,
INTERNODE_TRANSPORT_BACKEND_GRPC,
);
crate::rpc::runtime_sources::record_remote_disk_grpc_read_all_error();
Error::other(format!("can not get client, err: {err}"))
})?;
let request = Request::new(ReadAllRequest {
@@ -2188,34 +2157,21 @@ impl DiskAPI for RemoteDisk {
path: path.to_string(),
});
global_internode_metrics().record_outgoing_request_for_operation_and_backend(
INTERNODE_OPERATION_GRPC_READ_ALL,
INTERNODE_TRANSPORT_BACKEND_GRPC,
);
crate::rpc::runtime_sources::record_remote_disk_grpc_read_all_request();
let response = match client.read_all(request).await {
Ok(response) => response.into_inner(),
Err(err) => {
global_internode_metrics().record_error_for_operation_and_backend(
INTERNODE_OPERATION_GRPC_READ_ALL,
INTERNODE_TRANSPORT_BACKEND_GRPC,
);
crate::rpc::runtime_sources::record_remote_disk_grpc_read_all_error();
return Err(err.into());
}
};
if !response.success {
global_internode_metrics().record_error_for_operation_and_backend(
INTERNODE_OPERATION_GRPC_READ_ALL,
INTERNODE_TRANSPORT_BACKEND_GRPC,
);
crate::rpc::runtime_sources::record_remote_disk_grpc_read_all_error();
return Err(response.error.unwrap_or_default().into());
}
global_internode_metrics().record_recv_bytes_for_operation_and_backend(
INTERNODE_OPERATION_GRPC_READ_ALL,
INTERNODE_TRANSPORT_BACKEND_GRPC,
response.data.len(),
);
crate::rpc::runtime_sources::record_remote_disk_grpc_read_all_recv_bytes(response.data.len());
Ok(response.data)
},
get_max_timeout_duration(),
@@ -2972,7 +2928,7 @@ mod tests {
OpenWriteTestStep::Success,
]);
let remote_disk = new_remote_disk_with_transport(Arc::new(transport.clone())).await;
rustfs_io_metrics::internode_metrics::global_internode_metrics().reset_for_test();
crate::rpc::runtime_sources::reset_internode_metrics_for_test();
let _created = remote_disk
.create_file("orig-bucket", "bucket", "object/part.1", 4096)
@@ -2981,7 +2937,7 @@ mod tests {
let calls = transport.calls();
assert_eq!(calls.len(), 2, "create_file should retry exactly once");
let snapshot = rustfs_io_metrics::internode_metrics::global_internode_metrics().snapshot();
let snapshot = crate::rpc::runtime_sources::internode_metrics_snapshot_for_test();
assert_eq!(snapshot.outgoing_requests_total, 0);
}
+83
View File
@@ -0,0 +1,83 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use rustfs_io_metrics::internode_metrics::{
INTERNODE_OPERATION_GRPC_READ_ALL, INTERNODE_OPERATION_GRPC_WRITE_ALL, INTERNODE_OPERATION_PUT_FILE_STREAM,
INTERNODE_TRANSPORT_BACKEND_GRPC, INTERNODE_TRANSPORT_BACKEND_TCP_HTTP, global_internode_metrics,
};
#[cfg(test)]
use rustfs_io_metrics::internode_metrics::InternodeMetricsSnapshot;
pub(crate) fn record_remote_disk_open_write_retry(classification: &'static str) {
global_internode_metrics().record_retry_for_operation_and_backend(
INTERNODE_OPERATION_PUT_FILE_STREAM,
INTERNODE_TRANSPORT_BACKEND_TCP_HTTP,
classification,
);
}
pub(crate) fn record_remote_disk_open_write_retry_success(classification: &'static str) {
global_internode_metrics().record_retry_success_for_operation_and_backend(
INTERNODE_OPERATION_PUT_FILE_STREAM,
INTERNODE_TRANSPORT_BACKEND_TCP_HTTP,
classification,
);
}
pub(crate) fn record_remote_disk_grpc_write_all_error() {
global_internode_metrics()
.record_error_for_operation_and_backend(INTERNODE_OPERATION_GRPC_WRITE_ALL, INTERNODE_TRANSPORT_BACKEND_GRPC);
}
pub(crate) fn record_remote_disk_grpc_write_all_request() {
global_internode_metrics()
.record_outgoing_request_for_operation_and_backend(INTERNODE_OPERATION_GRPC_WRITE_ALL, INTERNODE_TRANSPORT_BACKEND_GRPC);
}
pub(crate) fn record_remote_disk_grpc_write_all_sent_bytes(bytes: usize) {
global_internode_metrics().record_sent_bytes_for_operation_and_backend(
INTERNODE_OPERATION_GRPC_WRITE_ALL,
INTERNODE_TRANSPORT_BACKEND_GRPC,
bytes,
);
}
pub(crate) fn record_remote_disk_grpc_read_all_error() {
global_internode_metrics()
.record_error_for_operation_and_backend(INTERNODE_OPERATION_GRPC_READ_ALL, INTERNODE_TRANSPORT_BACKEND_GRPC);
}
pub(crate) fn record_remote_disk_grpc_read_all_request() {
global_internode_metrics()
.record_outgoing_request_for_operation_and_backend(INTERNODE_OPERATION_GRPC_READ_ALL, INTERNODE_TRANSPORT_BACKEND_GRPC);
}
pub(crate) fn record_remote_disk_grpc_read_all_recv_bytes(bytes: usize) {
global_internode_metrics().record_recv_bytes_for_operation_and_backend(
INTERNODE_OPERATION_GRPC_READ_ALL,
INTERNODE_TRANSPORT_BACKEND_GRPC,
bytes,
);
}
#[cfg(test)]
pub(crate) fn reset_internode_metrics_for_test() {
global_internode_metrics().reset_for_test();
}
#[cfg(test)]
pub(crate) fn internode_metrics_snapshot_for_test() -> InternodeMetricsSnapshot {
global_internode_metrics().snapshot()
}
+5 -6
View File
@@ -16,11 +16,10 @@
// scoped to that module so generated internals do not relax lints elsewhere.
#[allow(unsafe_code)]
mod generated;
mod runtime_sources;
use proto_gen::node_service::node_service_client::NodeServiceClient;
use rustfs_common::{GLOBAL_CONN_MAP, evict_connection_with_log_level};
use rustfs_io_metrics::internode_metrics::global_internode_metrics;
use rustfs_tls_runtime::{load_global_outbound_tls_state, record_tls_consumer_stale_generation};
use std::{
collections::HashMap,
error::Error,
@@ -134,7 +133,7 @@ pub async fn create_new_channel(addr: &str) -> Result<Channel, Box<dyn Error>> {
// Overall timeout for any RPC - fail fast on unresponsive peers
.timeout(rpc_timeout);
let outbound_tls = load_global_outbound_tls_state().await;
let outbound_tls = runtime_sources::outbound_tls_state().await;
let generation = outbound_tls.generation.0;
let mut stale_generation = false;
{
@@ -180,11 +179,11 @@ pub async fn create_new_channel(addr: &str) -> Result<Channel, Box<dyn Error>> {
let channel = match connector.connect().await {
Ok(channel) => {
global_internode_metrics().record_dial_result(dial_started_at.elapsed(), true);
runtime_sources::record_grpc_dial_result(dial_started_at.elapsed(), true);
channel
}
Err(err) => {
global_internode_metrics().record_dial_result(dial_started_at.elapsed(), false);
runtime_sources::record_grpc_dial_result(dial_started_at.elapsed(), false);
return Err(err.into());
}
};
@@ -199,7 +198,7 @@ pub async fn create_new_channel(addr: &str) -> Result<Channel, Box<dyn Error>> {
generation_cache.insert(addr.to_string(), generation);
}
if stale_generation {
record_tls_consumer_stale_generation("protos_grpc_channel");
runtime_sources::record_stale_grpc_channel_tls_generation();
}
debug!("Successfully created and cached gRPC channel to: {}", addr);
+31
View File
@@ -0,0 +1,31 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use rustfs_io_metrics::internode_metrics::global_internode_metrics;
use rustfs_tls_runtime::{GlobalPublishedOutboundTlsState, load_global_outbound_tls_state, record_tls_consumer_stale_generation};
use std::time::Duration;
const PROTOS_GRPC_CHANNEL_TLS_CONSUMER: &str = "protos_grpc_channel";
pub(crate) async fn outbound_tls_state() -> GlobalPublishedOutboundTlsState {
load_global_outbound_tls_state().await
}
pub(crate) fn record_stale_grpc_channel_tls_generation() {
record_tls_consumer_stale_generation(PROTOS_GRPC_CHANNEL_TLS_CONSUMER);
}
pub(crate) fn record_grpc_dial_result(duration: Duration, success: bool) {
global_internode_metrics().record_dial_result(duration, success);
}