feat(connect): capture live RPC activity (#7807)

This commit is contained in:
Chris
2026-09-14 10:21:09 +08:00
committed by GitHub
parent 1cd9d1ed5d
commit b2420a0761
3 changed files with 216 additions and 21 deletions
+68 -6
View File
@@ -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)
}
+14 -5
View File
@@ -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]
+134 -10
View File
@@ -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::<Vec<_>>(),
["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());
}