From b2420a07610efef77ac4b8a5a1fd690982249a8b Mon Sep 17 00:00:00 2001 From: Chris Date: Mon, 14 Sep 2026 10:21:09 +0800 Subject: [PATCH] feat(connect): capture live RPC activity (#7807) --- rustfs/src/connect/diagnostics/top_rpc.rs | 74 ++++++++++- rustfs/src/storage/rpc/http_service.rs | 19 ++- rustfs/tests/connect_top_rpc.rs | 144 ++++++++++++++++++++-- 3 files changed, 216 insertions(+), 21 deletions(-) diff --git a/rustfs/src/connect/diagnostics/top_rpc.rs b/rustfs/src/connect/diagnostics/top_rpc.rs index d0978eef9..81cc4cde0 100644 --- a/rustfs/src/connect/diagnostics/top_rpc.rs +++ b/rustfs/src/connect/diagnostics/top_rpc.rs @@ -12,12 +12,15 @@ // See the License for the specific language governing permissions and // limitations under the License. -//! Honest `top.rpc` capability result. +//! Bounded `top.rpc` aggregation over pre-classified internode HTTP RPC completions. +use std::time::Duration; + +use rustfs_common::trace_bus::{TelemetryTraceOperation, TelemetryTraceStatus, subscribe_telemetry_trace_events}; use serde::Serialize; use tokio_util::sync::CancellationToken; -use super::top_api::{TopCaptureError, TopCaptureRequest, TopReasonCode, TopResult}; +use super::top_api::{MAX_SAFE_INTEGER, TopCaptureError, TopCaptureRequest, TopReasonCode, TopResult}; const TOOL_ID: &str = "top.rpc"; @@ -38,9 +41,68 @@ pub async fn capture_top_rpc( if cancel.is_cancelled() { return request.cancelled(TOOL_ID); } + let Some(_permit) = request.acquire(cancel).await? else { + return request.cancelled(TOOL_ID); + }; + // The typed source is emitted only after an internode HTTP response has a + // final status. Tonic streams need a separate body-completion boundary. + let mut subscription = subscribe_telemetry_trace_events(); + let started = tokio::time::Instant::now(); + let deadline = started + request.window; + let mut request_count = 0_u64; + let mut error_count = 0_u64; + let mut total_duration_micros = 0_u64; - // Internode metrics expose traffic, dial failures, and average dial time. - // They do not expose one matching RPC request/error/duration cohort, so v1 - // cannot be produced without mixing unrelated counters. - request.unsupported(TOOL_ID, TopReasonCode::UnsupportedTool) + loop { + tokio::select! { + biased; + () = cancel.cancelled() => return request.cancelled(TOOL_ID), + () = tokio::time::sleep_until(deadline) => break, + received = subscription.recv() => match received { + Ok(event) if event.operation == TelemetryTraceOperation::InternalRpc => { + if request_count >= u64::from(request.limits.max_operations) + || request_count >= u64::from(request.limits.max_records) + { + return request.failed(TOOL_ID, elapsed_millis(started.elapsed()), TopReasonCode::LimitExceeded); + } + let Ok(duration_micros) = u64::try_from(event.duration.as_micros()) else { + return request.failed(TOOL_ID, elapsed_millis(started.elapsed()), TopReasonCode::CollectionFailed); + }; + let Some(next_duration) = total_duration_micros.checked_add(duration_micros) else { + return request.failed(TOOL_ID, elapsed_millis(started.elapsed()), TopReasonCode::CollectionFailed); + }; + request_count += 1; + error_count += u64::from(event.status == TelemetryTraceStatus::Error); + total_duration_micros = next_duration; + } + Ok(_) => {} + Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => { + return request.failed(TOOL_ID, elapsed_millis(started.elapsed()), TopReasonCode::LimitExceeded); + } + Err(tokio::sync::broadcast::error::RecvError::Closed) => { + return request.failed(TOOL_ID, elapsed_millis(started.elapsed()), TopReasonCode::SourceUnavailable); + } + } + } + } + + request.validate_scope(TOOL_ID)?; + if request_count > MAX_SAFE_INTEGER || error_count > MAX_SAFE_INTEGER || total_duration_micros > MAX_SAFE_INTEGER { + return request.failed(TOOL_ID, elapsed_millis(started.elapsed()), TopReasonCode::CollectionFailed); + } + let window_millis = u64::try_from(request.window.as_millis()).map_err(|_| TopCaptureError::Limits)?; + request.succeeded( + TOOL_ID, + window_millis, + TopRpcData { + request_count, + error_count, + window_millis, + total_duration_micros, + }, + ) +} + +fn elapsed_millis(duration: Duration) -> u64 { + u64::try_from(duration.as_millis()).unwrap_or(u64::MAX).max(1) } diff --git a/rustfs/src/storage/rpc/http_service.rs b/rustfs/src/storage/rpc/http_service.rs index 9bf656ee6..2e10cec68 100644 --- a/rustfs/src/storage/rpc/http_service.rs +++ b/rustfs/src/storage/rpc/http_service.rs @@ -1777,17 +1777,26 @@ mod tests { #[tokio::test] #[serial_test::serial] - async fn internode_rpc_telemetry_publishes_only_the_classified_outcome() { + async fn internode_rpc_telemetry_publishes_only_classified_completions() { let mut subscription = rustfs_common::trace_bus::subscribe_telemetry_trace_events(); - super::emit_internode_rpc_telemetry(std::time::Instant::now(), StatusCode::FORBIDDEN); + super::emit_internode_rpc_telemetry(std::time::Instant::now() - Duration::from_micros(7), StatusCode::FORBIDDEN); + super::emit_internode_rpc_telemetry(std::time::Instant::now() - Duration::from_micros(11), StatusCode::OK); - let event = tokio::time::timeout(Duration::from_secs(1), subscription.recv()) + let error = tokio::time::timeout(Duration::from_secs(1), subscription.recv()) .await .expect("telemetry event should arrive") .expect("telemetry source should remain open"); - assert_eq!(event.operation, rustfs_common::trace_bus::TelemetryTraceOperation::InternalRpc); - assert_eq!(event.status, rustfs_common::trace_bus::TelemetryTraceStatus::Error); + let success = tokio::time::timeout(Duration::from_secs(1), subscription.recv()) + .await + .expect("telemetry event should arrive") + .expect("telemetry source should remain open"); + assert_eq!(error.operation, rustfs_common::trace_bus::TelemetryTraceOperation::InternalRpc); + assert_eq!(error.status, rustfs_common::trace_bus::TelemetryTraceStatus::Error); + assert!(error.duration > Duration::ZERO); + assert_eq!(success.operation, rustfs_common::trace_bus::TelemetryTraceOperation::InternalRpc); + assert_eq!(success.status, rustfs_common::trace_bus::TelemetryTraceStatus::Ok); + assert!(success.duration > Duration::ZERO); } #[tokio::test] diff --git a/rustfs/tests/connect_top_rpc.rs b/rustfs/tests/connect_top_rpc.rs index b2b7fd2a0..251bac911 100644 --- a/rustfs/tests/connect_top_rpc.rs +++ b/rustfs/tests/connect_top_rpc.rs @@ -17,13 +17,16 @@ use std::time::Duration; use rustfs::connect::diagnostics::{ LocalTopConsent, TopCaptureLimits, TopCaptureRequest, TopCaptureScope, TopOutcome, TopReasonCode, capture_top_rpc, }; +use rustfs_common::trace_bus::{ + TelemetryTraceEvent, TelemetryTraceOperation, TelemetryTraceStatus, telemetry_trace_emit, telemetry_trace_subscriber_count, +}; +use serial_test::serial; use time::OffsetDateTime; use tokio_util::sync::CancellationToken; -#[tokio::test] -async fn top_rpc_reports_unsupported_instead_of_mixing_network_and_trace_counters() { +fn request(window: Duration) -> TopCaptureRequest { let now = OffsetDateTime::now_utc().unix_timestamp(); - let request = TopCaptureRequest { + TopCaptureRequest { scope: TopCaptureScope { organization_name: "organizations/019e3ae0-0000-7000-8000-000000000010".to_owned(), cluster_name: "organizations/019e3ae0-0000-7000-8000-000000000010/clusters/019e3ae0-0000-7000-8000-000000000011".to_owned(), @@ -43,14 +46,135 @@ async fn top_rpc_reports_unsupported_instead_of_mixing_network_and_trace_counter }, }, limits: TopCaptureLimits::default(), - window: Duration::from_millis(1), + window, export_validity: Duration::from_secs(300), - }; + } +} - let result = capture_top_rpc(&request, &CancellationToken::new()) - .await - .expect("structured unsupported result"); - assert_eq!(result.outcome, TopOutcome::Unsupported); - assert_eq!(result.reason_code, TopReasonCode::UnsupportedTool); +async fn wait_for_subscription(previous: usize) { + tokio::time::timeout(Duration::from_secs(2), async { + while telemetry_trace_subscriber_count() <= previous { + tokio::task::yield_now().await; + } + }) + .await + .expect("top.rpc should subscribe to the runtime source"); +} + +#[tokio::test] +#[serial] +async fn top_rpc_captures_classified_rpc_outcomes_only() { + let request = request(Duration::from_millis(25)); + let previous = telemetry_trace_subscriber_count(); + let capture = tokio::spawn(async move { capture_top_rpc(&request, &CancellationToken::new()).await }); + wait_for_subscription(previous).await; + + assert!(telemetry_trace_emit(|| TelemetryTraceEvent::new( + TelemetryTraceOperation::InternalRpc, + Duration::from_micros(7), + TelemetryTraceStatus::Ok, + ))); + assert!(telemetry_trace_emit(|| TelemetryTraceEvent::new( + TelemetryTraceOperation::InternalRpc, + Duration::from_micros(11), + TelemetryTraceStatus::Error, + ))); + assert!(telemetry_trace_emit(|| TelemetryTraceEvent::new( + TelemetryTraceOperation::GetObject, + Duration::from_micros(13), + TelemetryTraceStatus::Error, + ))); + + let result = capture.await.expect("capture task").expect("top.rpc result"); + assert_eq!(result.outcome, TopOutcome::Succeeded); + assert_eq!(result.reason_code, TopReasonCode::Complete); + let data = result.data.expect("successful capture data"); + assert_eq!(data.request_count, 2); + assert_eq!(data.error_count, 1); + assert_eq!(data.total_duration_micros, 18); + let encoded = serde_json::to_value(data).expect("serialize top.rpc data"); + assert_eq!( + encoded.as_object().expect("top.rpc object").keys().collect::>(), + ["errorCount", "requestCount", "totalDurationMicros", "windowMillis"] + ); +} + +#[tokio::test] +#[serial] +async fn top_rpc_cancellation_stops_without_exportable_data() { + let request = request(Duration::from_secs(1)); + let cancel = CancellationToken::new(); + let task_cancel = cancel.clone(); + let previous = telemetry_trace_subscriber_count(); + let capture = tokio::spawn(async move { capture_top_rpc(&request, &task_cancel).await }); + wait_for_subscription(previous).await; + cancel.cancel(); + + let result = capture.await.expect("capture task").expect("cancelled result"); + assert_eq!(result.outcome, TopOutcome::Cancelled); + assert_eq!(result.reason_code, TopReasonCode::Cancelled); + assert!(result.data.is_none()); +} + +#[tokio::test] +#[serial] +async fn top_rpc_refuses_to_publish_a_truncated_window() { + let mut request = request(Duration::from_secs(1)); + request.limits.max_operations = 1; + request.limits.max_records = 1; + let previous = telemetry_trace_subscriber_count(); + let capture = tokio::spawn(async move { capture_top_rpc(&request, &CancellationToken::new()).await }); + wait_for_subscription(previous).await; + for _ in 0..2 { + assert!(telemetry_trace_emit(|| TelemetryTraceEvent::new( + TelemetryTraceOperation::InternalRpc, + Duration::from_micros(1), + TelemetryTraceStatus::Ok, + ))); + } + + let result = capture.await.expect("capture task").expect("limited result"); + assert_eq!(result.outcome, TopOutcome::Failed); + assert_eq!(result.reason_code, TopReasonCode::LimitExceeded); + assert!(result.data.is_none()); +} + +#[tokio::test(flavor = "current_thread")] +#[serial] +async fn top_rpc_fails_closed_when_the_source_lags() { + let request = request(Duration::from_secs(1)); + let previous = telemetry_trace_subscriber_count(); + let capture = tokio::spawn(async move { capture_top_rpc(&request, &CancellationToken::new()).await }); + wait_for_subscription(previous).await; + for _ in 0..1_025 { + assert!(telemetry_trace_emit(|| TelemetryTraceEvent::new( + TelemetryTraceOperation::InternalRpc, + Duration::from_micros(1), + TelemetryTraceStatus::Ok, + ))); + } + + let result = capture.await.expect("capture task").expect("lagged result"); + assert_eq!(result.outcome, TopOutcome::Failed); + assert_eq!(result.reason_code, TopReasonCode::LimitExceeded); + assert!(result.data.is_none()); +} + +#[tokio::test] +#[serial] +async fn top_rpc_rejects_duration_overflow() { + let request = request(Duration::from_secs(1)); + let previous = telemetry_trace_subscriber_count(); + let capture = tokio::spawn(async move { capture_top_rpc(&request, &CancellationToken::new()).await }); + wait_for_subscription(previous).await; + assert!(telemetry_trace_emit(|| TelemetryTraceEvent::new( + TelemetryTraceOperation::InternalRpc, + Duration::MAX, + TelemetryTraceStatus::Ok, + ))); + + let result = capture.await.expect("capture task").expect("overflow result"); + assert_eq!(result.outcome, TopOutcome::Failed); + assert_eq!(result.reason_code, TopReasonCode::CollectionFailed); assert!(result.data.is_none()); }