diff --git a/quickwit/quickwit-common/src/tower/metrics.rs b/quickwit/quickwit-common/src/tower/metrics.rs index 17a9f6acec3..cb3cc9884e3 100644 --- a/quickwit/quickwit-common/src/tower/metrics.rs +++ b/quickwit/quickwit-common/src/tower/metrics.rs @@ -22,6 +22,7 @@ use quickwit_metrics::{ Counter, Gauge, Histogram, Labels, LazyCounter, LazyGauge, LazyHistogram, counter, gauge, histogram, labels, lazy_counter, lazy_gauge, lazy_histogram, }; +use tower::retry::Policy; use tower::{Layer, Service}; use crate::metrics::exponential_buckets; @@ -141,6 +142,37 @@ pub struct GrpcMetricsLayer { request_duration_seconds: Histogram, } +/// Retry policy that records each failed attempt which will be retried. +#[derive(Clone)] +pub struct GrpcRetryPolicy

{ + inner: P, + metrics_layer: GrpcMetricsLayer, +} + +impl Policy for GrpcRetryPolicy

+where + P: Policy, + R: RpcName, + E: GrpcStatusCode, +{ + type Future = P::Future; + + fn retry(&mut self, request: &mut R, result: &mut Result) -> Option { + let retry = self.inner.retry(request, result); + if let Err(error) = result + && retry.is_some() + { + self.metrics_layer + .record_request(R::rpc_name(), "retry", error.grpc_status_code()); + } + retry + } + + fn clone_request(&mut self, request: &R) -> Option { + self.inner.clone_request(request) + } +} + impl GrpcMetricsLayer { pub fn new(subsystem: &'static str, kind: &'static str) -> Self { let labels = Self::default_labels(subsystem, kind); @@ -164,6 +196,26 @@ impl GrpcMetricsLayer { } } + /// Wraps a retry policy to record failed attempts that it retries. + pub fn from_retry_policy

(&self, inner: P) -> GrpcRetryPolicy

{ + GrpcRetryPolicy { + inner, + metrics_layer: self.clone(), + } + } + + fn record_request(&self, rpc_name: &'static str, status: &'static str, code: tonic::Code) { + counter!( + parent: self.requests_total, + labels: [labels!( + "rpc" => rpc_name, + "status" => status, + "code" => grpc_code_label(code), + )], + ) + .inc(); + } + fn default_labels(subsystem: &'static str, kind: &'static str) -> Labels<3> { // `service` is kept for backward compatibility with existing consumers. Prefer // `grpc_service` for new consumers. @@ -243,9 +295,11 @@ mod tests { use metrics::with_local_recorder; use metrics_util::debugging::{DebugValue, DebuggingRecorder}; + use super::super::retry::{RetryLayer, RetryPolicy}; use super::*; + use crate::retry::{RetryParams, Retryable}; - #[derive(Debug)] + #[derive(Clone, Debug)] struct HelloRequest; impl RpcName for HelloRequest { @@ -262,13 +316,29 @@ mod tests { } } + #[derive(Debug)] + struct RetryableError; + + impl GrpcStatusCode for RetryableError { + fn grpc_status_code(&self) -> tonic::Code { + tonic::Code::Unavailable + } + } + + impl Retryable for RetryableError { + fn is_retryable(&self) -> bool { + true + } + } + #[test] fn test_grpc_metrics() { let recorder = DebuggingRecorder::new(); let snapshotter = recorder.snapshotter(); with_local_recorder(&recorder, || { - futures::executor::block_on(async { + let runtime = tokio::runtime::Runtime::new().unwrap(); + runtime.block_on(async { let primary_layer = GrpcMetricsLayer::new_with_labels( "quickwit_test", "server", @@ -289,10 +359,11 @@ mod tests { let mut read_replica_service = read_replica_layer.layer(tower::service_fn( |request: HelloRequest| async move { Ok::<_, tonic::Status>(request) }, )); - let mut failing_service = - primary_layer.layer(tower::service_fn(|_request: HelloRequest| async move { + let mut failing_service = primary_layer.clone().layer(tower::service_fn( + |_request: HelloRequest| async move { Err::(tonic::Status::not_found("not found")) - })); + }, + )); hello_service.call(HelloRequest).await.unwrap(); goodbye_service.call(GoodbyeRequest).await.unwrap(); @@ -301,6 +372,19 @@ mod tests { let hello_future = hello_service.call(HelloRequest); drop(hello_future); + + let retry_params = RetryParams { + base_delay: std::time::Duration::ZERO, + max_delay: std::time::Duration::ZERO, + max_attempts: 3, + }; + let retry_policy = primary_layer.from_retry_policy(RetryPolicy::from(retry_params)); + let mut retry_service = primary_layer.layer(RetryLayer::new(retry_policy).layer( + tower::service_fn(|_request: HelloRequest| async move { + Err::(RetryableError) + }), + )); + retry_service.call(HelloRequest).await.unwrap_err(); }); }); @@ -347,5 +431,13 @@ mod tests { counter_value("hello", "error", "not_found", "primary"), Some(&DebugValue::Counter(1)) ); + assert_eq!( + counter_value("hello", "retry", "unavailable", "primary"), + Some(&DebugValue::Counter(2)) + ); + assert_eq!( + counter_value("hello", "error", "unavailable", "primary"), + Some(&DebugValue::Counter(1)) + ); } } diff --git a/quickwit/quickwit-common/src/tower/mod.rs b/quickwit/quickwit-common/src/tower/mod.rs index c763ce22bf8..67748cf766b 100644 --- a/quickwit/quickwit-common/src/tower/mod.rs +++ b/quickwit/quickwit-common/src/tower/mod.rs @@ -44,7 +44,7 @@ pub use estimate_rate::{EstimateRate, EstimateRateLayer}; pub use event_listener::{EventListener, EventListenerLayer}; use futures::Future; pub use load_shed::{LoadShed, LoadShedLayer, MakeLoadShedError}; -pub use metrics::{GrpcMetrics, GrpcMetricsLayer, GrpcStatusCode, RpcName}; +pub use metrics::{GrpcMetrics, GrpcMetricsLayer, GrpcRetryPolicy, GrpcStatusCode, RpcName}; pub use one_task_per_call_layer::{OneTaskPerCallLayer, TaskCancelled}; pub use pool::Pool; pub use rate::{ConstantRate, Rate}; diff --git a/quickwit/quickwit-serve/src/metastore.rs b/quickwit/quickwit-serve/src/metastore.rs index a4ee434232b..f42dc80dbf5 100644 --- a/quickwit/quickwit-serve/src/metastore.rs +++ b/quickwit/quickwit-serve/src/metastore.rs @@ -176,10 +176,13 @@ impl LocalMetastoreServer { } _ => unreachable!("unexpected metastore service `{service}`"), }; + let retry_policy = + metrics_layer.from_retry_policy(RetryPolicy::from(RetryParams::standard())); Ok(MetastoreServiceClient::tower() - .stack_layer(RetryLayer::new(RetryPolicy::from(RetryParams::standard()))) - .stack_layer(TimeoutLayer::new(GRPC_METASTORE_SERVICE_TIMEOUT)) + // Metrics wrap retries so they record only the final outcome. .stack_layer(metrics_layer) + .stack_layer(RetryLayer::new(retry_policy)) + .stack_layer(TimeoutLayer::new(GRPC_METASTORE_SERVICE_TIMEOUT)) .stack_layer(tower::limit::GlobalConcurrencyLimitLayer::new( get_metastore_client_max_concurrency(), ))