From ce6dd9d37140e5e2c630032ac1416d99332c9bbb Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Mon, 17 Aug 2026 09:55:15 +0200 Subject: [PATCH 01/12] Add metastore retry outcome metrics --- docs/reference/metrics.md | 19 +++-- quickwit/quickwit-common/src/tower/metrics.rs | 10 ++- quickwit/quickwit-common/src/tower/retry.rs | 77 ++++++++++++++++++- quickwit/quickwit-serve/src/metastore.rs | 9 ++- 4 files changed, 104 insertions(+), 11 deletions(-) diff --git a/docs/reference/metrics.md b/docs/reference/metrics.md index 33a49854895..d28984e8745 100644 --- a/docs/reference/metrics.md +++ b/docs/reference/metrics.md @@ -50,15 +50,24 @@ Quickwit exposes metrics for several cache components, including `fastfields`, ` ## Metastore Metrics -All metastore methods are monitored by the 3 metrics: +Metastore RPCs use the shared gRPC metrics: | Namespace | Metric Name | Description | Labels | Type | | --------- | ----------- | ----------- | ------ | ---- | -| `quickwit_metastore` | `requests_total` | Number of requests | [`operation`, `index`] | `counter` | -| `quickwit_metastore` | `request_errors_total` | Number of failed requests | [`operation`, `index`] | `counter` | -| `quickwit_metastore` | `request_duration_seconds` | Duration of requests | [`operation`, `index`, `error`] | `histogram` | +| `quickwit_grpc` | `requests_total` | Number of metastore RPC outcomes and retryable failed attempts | [`grpc_service`, `kind`, `metastore_kind`, `rpc`, `status`, `code`] | `counter` | -Examples of operation names: `create_index`, `index_metadata`, `delete_index`, `stage_splits`, `publish_splits`, `list_splits`, `add_source`, ... +Filter this metric with `grpc_service="metastore"`. `kind` identifies the client or server, +and `metastore_kind` identifies the primary or read-replica metastore. The `status` label has the +following meaning for client metrics: + +- `success`: the logical RPC succeeded, including after any retries. +- `transient`: a retryable failed attempt for which the client will issue another attempt. +- `error`: the logical RPC failed and was returned to the caller (a non-retryable error or after + retry exhaustion). +- `cancelled`: the client cancelled the RPC before it completed. + +`rpc` contains the metastore operation name, such as `create_index`, `index_metadata`, +`delete_index`, `stage_splits`, `publish_splits`, `list_splits`, or `add_source`. PostgreSQL-backed metastores also expose connection pool gauges: diff --git a/quickwit/quickwit-common/src/tower/metrics.rs b/quickwit/quickwit-common/src/tower/metrics.rs index 17a9f6acec3..71db523e9ed 100644 --- a/quickwit/quickwit-common/src/tower/metrics.rs +++ b/quickwit/quickwit-common/src/tower/metrics.rs @@ -47,7 +47,7 @@ impl GrpcStatusCode for std::convert::Infallible { } } -fn grpc_code_label(code: tonic::Code) -> &'static str { +pub(crate) fn grpc_code_label(code: tonic::Code) -> &'static str { match code { tonic::Code::Ok => "ok", tonic::Code::Cancelled => "cancelled", @@ -164,6 +164,14 @@ impl GrpcMetricsLayer { } } + /// Returns the request counter with this layer's static labels. + /// + /// A retry policy uses this to record a retryable failed attempt with + /// `status="transient"` after it decides to retry it. + pub fn requests_total_counter(&self) -> Counter { + self.requests_total.clone() + } + 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. diff --git a/quickwit/quickwit-common/src/tower/retry.rs b/quickwit/quickwit-common/src/tower/retry.rs index 4594aaa1cda..a07bd2d0701 100644 --- a/quickwit/quickwit-common/src/tower/retry.rs +++ b/quickwit/quickwit-common/src/tower/retry.rs @@ -15,11 +15,13 @@ use std::any::type_name; use std::fmt; +use quickwit_metrics::{Counter, counter, labels}; use tokio::time::Sleep; use tower::Layer; use tower::retry::{Policy, Retry}; use tracing::debug; +use super::metrics::{GrpcStatusCode, RpcName, grpc_code_label}; use crate::retry::{RetryParams, Retryable}; /// Retry layer copy/pasted from `tower::retry::RetryLayer` @@ -47,10 +49,19 @@ impl

RetryLayer

{ } } -#[derive(Clone, Copy, Debug)] +#[derive(Clone, Debug)] pub struct RetryPolicy { num_attempts: usize, retry_params: RetryParams, + retry_metrics_counter_opt: Option, +} + +impl RetryPolicy { + /// Records each failed attempt that this policy will retry as a transient request. + pub fn with_retry_metrics(mut self, retry_metrics_counter: Counter) -> Self { + self.retry_metrics_counter_opt = Some(retry_metrics_counter); + self + } } impl From for RetryPolicy { @@ -58,14 +69,15 @@ impl From for RetryPolicy { Self { num_attempts: 0, retry_params, + retry_metrics_counter_opt: None, } } } impl Policy for RetryPolicy where - R: Clone, - E: fmt::Debug + Retryable, + R: Clone + RpcName, + E: fmt::Debug + Retryable + GrpcStatusCode, { type Future = Sleep; @@ -78,6 +90,17 @@ where if !error.is_retryable() || self.num_attempts >= self.retry_params.max_attempts { None } else { + if let Some(retry_metrics_counter) = &self.retry_metrics_counter_opt { + counter!( + parent: retry_metrics_counter, + labels: [labels!( + "rpc" => R::rpc_name(), + "status" => "transient", + "code" => grpc_code_label(error.grpc_status_code()), + )], + ) + .inc(); + } let delay = self.retry_params.compute_delay(self.num_attempts); debug!( num_attempts=%self.num_attempts, @@ -104,6 +127,8 @@ mod tests { use std::task::{Context, Poll}; use futures::future::{Ready, ready}; + use metrics::with_local_recorder; + use metrics_util::debugging::{DebugValue, DebuggingRecorder}; use tower::{Layer, Service, ServiceExt}; use super::*; @@ -123,6 +148,12 @@ mod tests { } } + impl GrpcStatusCode for Retry { + fn grpc_status_code(&self) -> tonic::Code { + tonic::Code::Unavailable + } + } + #[derive(Debug, Clone, Default)] struct HelloService; @@ -134,6 +165,12 @@ mod tests { results: HelloResults, } + impl RpcName for HelloRequest { + fn rpc_name() -> &'static str { + "hello" + } + } + impl Service for HelloService { type Response = (); type Error = Retry<()>; @@ -155,6 +192,40 @@ mod tests { } } + #[tokio::test] + async fn test_retry_policy_records_retryable_failures_as_transient() { + let recorder = DebuggingRecorder::new(); + let snapshotter = recorder.snapshotter(); + + with_local_recorder(&recorder, || { + let retry_metrics_counter = counter!( + name: "requests_total", + description: "test request count", + subsystem: "grpc", + ); + let mut retry_policy = RetryPolicy::from(RetryParams::for_test()) + .with_retry_metrics(retry_metrics_counter); + let mut request = HelloRequest::default(); + let mut result: Result<(), Retry<()>> = Err(Retry::Transient(())); + + assert!(retry_policy.retry(&mut request, &mut result).is_some()); + }); + + let snapshot = snapshotter.snapshot().into_vec(); + assert!(snapshot.iter().any(|(composite_key, _, _, value)| { + let (_, key) = composite_key.clone().into_parts(); + let labels = key + .labels() + .map(|label| (label.key(), label.value())) + .collect::>(); + key.name() == "quickwit_grpc_requests_total" + && labels.contains(&("rpc", "hello")) + && labels.contains(&("status", "transient")) + && labels.contains(&("code", "unavailable")) + && value == &DebugValue::Counter(1) + })); + } + #[tokio::test] async fn test_retry_policy() { let retry_policy = RetryPolicy::from(RetryParams::for_test()); diff --git a/quickwit/quickwit-serve/src/metastore.rs b/quickwit/quickwit-serve/src/metastore.rs index a4ee434232b..5b9a693e1c6 100644 --- a/quickwit/quickwit-serve/src/metastore.rs +++ b/quickwit/quickwit-serve/src/metastore.rs @@ -176,10 +176,15 @@ impl LocalMetastoreServer { } _ => unreachable!("unexpected metastore service `{service}`"), }; + let retry_policy = RetryPolicy::from(RetryParams::standard()) + .with_retry_metrics(metrics_layer.requests_total_counter()); Ok(MetastoreServiceClient::tower() - .stack_layer(RetryLayer::new(RetryPolicy::from(RetryParams::standard()))) - .stack_layer(TimeoutLayer::new(GRPC_METASTORE_SERVICE_TIMEOUT)) + // Keep the request metric outside retries so `success` and `error` represent the + // final result seen by the caller. RetryPolicy records each retryable failed attempt + // with `status="transient"`. .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(), )) From 4915a1300fff8e37485c03f171848954f7472b5f Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Mon, 17 Aug 2026 11:42:08 +0200 Subject: [PATCH 02/12] Decouple retry metrics with callbacks --- quickwit/quickwit-common/src/tower/metrics.rs | 23 ++- quickwit/quickwit-common/src/tower/mod.rs | 2 +- quickwit/quickwit-common/src/tower/retry.rs | 177 +++++++++++------- quickwit/quickwit-serve/src/metastore.rs | 26 ++- 4 files changed, 146 insertions(+), 82 deletions(-) diff --git a/quickwit/quickwit-common/src/tower/metrics.rs b/quickwit/quickwit-common/src/tower/metrics.rs index 71db523e9ed..8ef18cbd421 100644 --- a/quickwit/quickwit-common/src/tower/metrics.rs +++ b/quickwit/quickwit-common/src/tower/metrics.rs @@ -164,12 +164,17 @@ impl GrpcMetricsLayer { } } - /// Returns the request counter with this layer's static labels. - /// - /// A retry policy uses this to record a retryable failed attempt with - /// `status="transient"` after it decides to retry it. - pub fn requests_total_counter(&self) -> Counter { - self.requests_total.clone() + /// Records a request outcome with this layer's static labels. + pub 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> { @@ -288,6 +293,8 @@ mod tests { labels!("metastore_kind" => "read_replica", "test_label" => "test"), ); + primary_layer.record_request("hello", "transient", tonic::Code::Unavailable); + let mut hello_service = primary_layer.clone().layer(tower::service_fn( |request: HelloRequest| async move { Ok::<_, tonic::Status>(request) }, )); @@ -339,6 +346,10 @@ mod tests { counter_value("hello", "success", "ok", "primary"), Some(&DebugValue::Counter(1)) ); + assert_eq!( + counter_value("hello", "transient", "unavailable", "primary"), + Some(&DebugValue::Counter(1)) + ); assert_eq!( counter_value("goodbye", "success", "ok", "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..d243d556e70 100644 --- a/quickwit/quickwit-common/src/tower/mod.rs +++ b/quickwit/quickwit-common/src/tower/mod.rs @@ -50,7 +50,7 @@ pub use pool::Pool; pub use rate::{ConstantRate, Rate}; pub use rate_estimator::{RateEstimator, SmaRateEstimator}; pub use rate_limit::{RateLimit, RateLimitLayer}; -pub use retry::{RetryLayer, RetryPolicy}; +pub use retry::{RetryCallbacks, RetryLayer, RetryPolicy}; pub use timeout::{Timeout, TimeoutExceeded, TimeoutLayer}; pub use transport::{BalanceChannel, warmup_channel}; diff --git a/quickwit/quickwit-common/src/tower/retry.rs b/quickwit/quickwit-common/src/tower/retry.rs index a07bd2d0701..7ed2435fb05 100644 --- a/quickwit/quickwit-common/src/tower/retry.rs +++ b/quickwit/quickwit-common/src/tower/retry.rs @@ -15,19 +15,18 @@ use std::any::type_name; use std::fmt; -use quickwit_metrics::{Counter, counter, labels}; use tokio::time::Sleep; use tower::Layer; use tower::retry::{Policy, Retry}; use tracing::debug; -use super::metrics::{GrpcStatusCode, RpcName, grpc_code_label}; use crate::retry::{RetryParams, Retryable}; /// Retry layer copy/pasted from `tower::retry::RetryLayer` /// but which implements `Clone`. impl Layer for RetryLayer

-where P: Clone +where + P: Clone, { type Service = Retry; @@ -49,18 +48,39 @@ impl

RetryLayer

{ } } +/// Callbacks invoked after the retry policy classifies an attempt's result. +pub trait RetryCallbacks: Clone { + /// Called when an attempt failed and another attempt will be made. + fn on_retry(&mut self, _request: &R, _error: &E, _num_attempts: usize) {} + + /// Called when the logical request completed with an error. + fn on_error(&mut self, _request: &R, _error: &E, _num_attempts: usize) {} + + /// Called when the logical request completed successfully. + fn on_success(&mut self, _request: &R, _response: &T, _num_attempts: usize) {} +} + +/// Default callbacks that do nothing. +#[derive(Clone, Copy, Debug, Default)] +pub struct NoopRetryCallbacks; + +impl RetryCallbacks for NoopRetryCallbacks {} + #[derive(Clone, Debug)] -pub struct RetryPolicy { +pub struct RetryPolicy { num_attempts: usize, retry_params: RetryParams, - retry_metrics_counter_opt: Option, + callbacks: C, } -impl RetryPolicy { - /// Records each failed attempt that this policy will retry as a transient request. - pub fn with_retry_metrics(mut self, retry_metrics_counter: Counter) -> Self { - self.retry_metrics_counter_opt = Some(retry_metrics_counter); - self +impl RetryPolicy { + /// Replaces the callbacks invoked when an attempt is classified. + pub fn with_callbacks(self, callbacks: D) -> RetryPolicy { + RetryPolicy { + num_attempts: self.num_attempts, + retry_params: self.retry_params, + callbacks, + } } } @@ -69,38 +89,34 @@ impl From for RetryPolicy { Self { num_attempts: 0, retry_params, - retry_metrics_counter_opt: None, + callbacks: NoopRetryCallbacks, } } } -impl Policy for RetryPolicy +impl Policy for RetryPolicy where - R: Clone + RpcName, - E: fmt::Debug + Retryable + GrpcStatusCode, + R: Clone, + E: fmt::Debug + Retryable, + C: RetryCallbacks, { type Future = Sleep; - fn retry(&mut self, _request: &mut R, result: &mut Result) -> Option { + fn retry(&mut self, request: &mut R, result: &mut Result) -> Option { match result { - Ok(_) => None, + Ok(response) => { + self.callbacks + .on_success(request, response, self.num_attempts + 1); + None + } Err(error) => { self.num_attempts += 1; if !error.is_retryable() || self.num_attempts >= self.retry_params.max_attempts { + self.callbacks.on_error(request, error, self.num_attempts); None } else { - if let Some(retry_metrics_counter) = &self.retry_metrics_counter_opt { - counter!( - parent: retry_metrics_counter, - labels: [labels!( - "rpc" => R::rpc_name(), - "status" => "transient", - "code" => grpc_code_label(error.grpc_status_code()), - )], - ) - .inc(); - } + self.callbacks.on_retry(request, error, self.num_attempts); let delay = self.retry_params.compute_delay(self.num_attempts); debug!( num_attempts=%self.num_attempts, @@ -127,8 +143,6 @@ mod tests { use std::task::{Context, Poll}; use futures::future::{Ready, ready}; - use metrics::with_local_recorder; - use metrics_util::debugging::{DebugValue, DebuggingRecorder}; use tower::{Layer, Service, ServiceExt}; use super::*; @@ -148,12 +162,6 @@ mod tests { } } - impl GrpcStatusCode for Retry { - fn grpc_status_code(&self) -> tonic::Code { - tonic::Code::Unavailable - } - } - #[derive(Debug, Clone, Default)] struct HelloService; @@ -165,12 +173,6 @@ mod tests { results: HelloResults, } - impl RpcName for HelloRequest { - fn rpc_name() -> &'static str { - "hello" - } - } - impl Service for HelloService { type Response = (); type Error = Retry<()>; @@ -192,38 +194,69 @@ mod tests { } } + #[derive(Clone, Debug, Default)] + struct TestRetryCallbacks { + events: Arc>>, + } + + impl RetryCallbacks> for TestRetryCallbacks { + fn on_retry(&mut self, _request: &HelloRequest, _error: &Retry<()>, num_attempts: usize) { + self.events.lock().unwrap().push(("retry", num_attempts)); + } + + fn on_error(&mut self, _request: &HelloRequest, _error: &Retry<()>, num_attempts: usize) { + self.events.lock().unwrap().push(("error", num_attempts)); + } + + fn on_success(&mut self, _request: &HelloRequest, _response: &(), num_attempts: usize) { + self.events.lock().unwrap().push(("success", num_attempts)); + } + } + #[tokio::test] - async fn test_retry_policy_records_retryable_failures_as_transient() { - let recorder = DebuggingRecorder::new(); - let snapshotter = recorder.snapshotter(); - - with_local_recorder(&recorder, || { - let retry_metrics_counter = counter!( - name: "requests_total", - description: "test request count", - subsystem: "grpc", - ); - let mut retry_policy = RetryPolicy::from(RetryParams::for_test()) - .with_retry_metrics(retry_metrics_counter); - let mut request = HelloRequest::default(); - let mut result: Result<(), Retry<()>> = Err(Retry::Transient(())); - - assert!(retry_policy.retry(&mut request, &mut result).is_some()); - }); - - let snapshot = snapshotter.snapshot().into_vec(); - assert!(snapshot.iter().any(|(composite_key, _, _, value)| { - let (_, key) = composite_key.clone().into_parts(); - let labels = key - .labels() - .map(|label| (label.key(), label.value())) - .collect::>(); - key.name() == "quickwit_grpc_requests_total" - && labels.contains(&("rpc", "hello")) - && labels.contains(&("status", "transient")) - && labels.contains(&("code", "unavailable")) - && value == &DebugValue::Counter(1) - })); + async fn test_retry_policy_callbacks() { + let callbacks = TestRetryCallbacks::default(); + let mut retry_policy = + RetryPolicy::from(RetryParams::for_test()).with_callbacks(callbacks.clone()); + let mut request = HelloRequest::default(); + + let mut transient_error = Err(Retry::Transient(())); + assert!( + retry_policy + .retry(&mut request, &mut transient_error) + .is_some() + ); + let mut success = Ok(()); + assert!(retry_policy.retry(&mut request, &mut success).is_none()); + assert_eq!( + callbacks.events.lock().unwrap().as_slice(), + &[("retry", 1), ("success", 2)] + ); + + let callbacks = TestRetryCallbacks::default(); + let mut retry_policy = + RetryPolicy::from(RetryParams::for_test()).with_callbacks(callbacks.clone()); + for num_attempts in 1..=3 { + let mut transient_error = Err(Retry::Transient(())); + let retry = retry_policy.retry(&mut request, &mut transient_error); + assert_eq!(retry.is_some(), num_attempts < 3); + } + assert_eq!( + callbacks.events.lock().unwrap().as_slice(), + &[("retry", 1), ("retry", 2), ("error", 3)] + ); + + let callbacks = TestRetryCallbacks::default(); + let mut retry_policy = + RetryPolicy::from(RetryParams::for_test()).with_callbacks(callbacks.clone()); + let mut permanent_error = Err(Retry::Permanent(())); + assert!(retry_policy + .retry(&mut request, &mut permanent_error) + .is_none()); + assert_eq!( + callbacks.events.lock().unwrap().as_slice(), + &[("error", 1)] + ); } #[tokio::test] diff --git a/quickwit/quickwit-serve/src/metastore.rs b/quickwit/quickwit-serve/src/metastore.rs index 5b9a693e1c6..9999a90593b 100644 --- a/quickwit/quickwit-serve/src/metastore.rs +++ b/quickwit/quickwit-serve/src/metastore.rs @@ -21,7 +21,8 @@ use quickwit_cluster::Cluster; use quickwit_common::pubsub::EventBroker; use quickwit_common::retry::RetryParams; use quickwit_common::tower::{ - EventListenerLayer, GrpcMetricsLayer, LoadShedLayer, RetryLayer, RetryPolicy, TimeoutLayer, + EventListenerLayer, GrpcMetricsLayer, GrpcStatusCode, LoadShedLayer, RetryCallbacks, + RetryLayer, RetryPolicy, RpcName, TimeoutLayer, }; use quickwit_common::uri::Uri; use quickwit_config::service::QuickwitService; @@ -71,6 +72,22 @@ static READ_REPLICA_METASTORE_GRPC_SERVER_METRICS_LAYER: LazyLock RetryCallbacks for MetastoreRetryCallbacks +where + R: RpcName, + E: GrpcStatusCode, +{ + fn on_retry(&mut self, _request: &R, error: &E, _num_attempts: usize) { + self.metrics_layer + .record_request(R::rpc_name(), "transient", error.grpc_status_code()); + } +} + fn get_metastore_client_max_concurrency() -> usize { quickwit_common::get_from_env( METASTORE_CLIENT_MAX_CONCURRENCY_ENV_KEY, @@ -176,8 +193,11 @@ impl LocalMetastoreServer { } _ => unreachable!("unexpected metastore service `{service}`"), }; - let retry_policy = RetryPolicy::from(RetryParams::standard()) - .with_retry_metrics(metrics_layer.requests_total_counter()); + let retry_callbacks = MetastoreRetryCallbacks { + metrics_layer: metrics_layer.clone(), + }; + let retry_policy = + RetryPolicy::from(RetryParams::standard()).with_callbacks(retry_callbacks); Ok(MetastoreServiceClient::tower() // Keep the request metric outside retries so `success` and `error` represent the // final result seen by the caller. RetryPolicy records each retryable failed attempt From 4096b7a921f3f7445e2a32888756c91a11e20f10 Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Mon, 17 Aug 2026 11:55:02 +0200 Subject: [PATCH 03/12] Simplify gRPC retry metrics --- quickwit/quickwit-common/src/tower/metrics.rs | 79 ++++++++++-- quickwit/quickwit-common/src/tower/mod.rs | 4 +- quickwit/quickwit-common/src/tower/retry.rs | 116 +----------------- quickwit/quickwit-serve/src/metastore.rs | 30 +---- 4 files changed, 83 insertions(+), 146 deletions(-) diff --git a/quickwit/quickwit-common/src/tower/metrics.rs b/quickwit/quickwit-common/src/tower/metrics.rs index 8ef18cbd421..3d920f33379 100644 --- a/quickwit/quickwit-common/src/tower/metrics.rs +++ b/quickwit/quickwit-common/src/tower/metrics.rs @@ -12,6 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. +use std::fmt; use std::pin::Pin; use std::task::{Context, Poll}; use std::time::Instant; @@ -22,9 +23,12 @@ 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 super::retry::RetryPolicy; use crate::metrics::exponential_buckets; +use crate::retry::{RetryParams, Retryable}; pub trait RpcName { fn rpc_name() -> &'static str; @@ -47,7 +51,7 @@ impl GrpcStatusCode for std::convert::Infallible { } } -pub(crate) fn grpc_code_label(code: tonic::Code) -> &'static str { +fn grpc_code_label(code: tonic::Code) -> &'static str { match code { tonic::Code::Ok => "ok", tonic::Code::Cancelled => "cancelled", @@ -141,6 +145,37 @@ pub struct GrpcMetricsLayer { request_duration_seconds: Histogram, } +/// Retry policy that records each failed attempt which will be retried as transient. +#[derive(Clone)] +pub struct GrpcRetryPolicy { + inner: RetryPolicy, + metrics_layer: GrpcMetricsLayer, +} + +impl Policy for GrpcRetryPolicy +where + R: Clone + RpcName, + E: fmt::Debug + Retryable + GrpcStatusCode, +{ + type Future = >::Future; + + fn retry(&mut self, request: &mut R, result: &mut Result) -> Option { + let retry = self.inner.retry(request, result); + if retry.is_some() { + let Err(error) = result else { + unreachable!("a successful request should not be retried"); + }; + self.metrics_layer + .record_request(R::rpc_name(), "transient", error.grpc_status_code()); + } + retry + } + + fn clone_request(&mut self, request: &R) -> Option { + >::clone_request(&mut self.inner, request) + } +} + impl GrpcMetricsLayer { pub fn new(subsystem: &'static str, kind: &'static str) -> Self { let labels = Self::default_labels(subsystem, kind); @@ -164,8 +199,15 @@ impl GrpcMetricsLayer { } } - /// Records a request outcome with this layer's static labels. - pub fn record_request(&self, rpc_name: &'static str, status: &'static str, code: tonic::Code) { + /// Wraps the retry policy to record retryable failed attempts as transient requests. + pub fn retry_policy(&self, retry_params: RetryParams) -> GrpcRetryPolicy { + GrpcRetryPolicy { + inner: RetryPolicy::from(retry_params), + 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!( @@ -258,7 +300,7 @@ mod tests { use super::*; - #[derive(Debug)] + #[derive(Clone, Debug)] struct HelloRequest; impl RpcName for HelloRequest { @@ -269,14 +311,29 @@ mod tests { struct GoodbyeRequest; + #[derive(Debug)] + struct RetryableError; + + impl Retryable for RetryableError { + fn is_retryable(&self) -> bool { + true + } + } + + impl GrpcStatusCode for RetryableError { + fn grpc_status_code(&self) -> tonic::Code { + tonic::Code::Unavailable + } + } + impl RpcName for GoodbyeRequest { fn rpc_name() -> &'static str { "goodbye" } } - #[test] - fn test_grpc_metrics() { + #[tokio::test] + async fn test_grpc_metrics() { let recorder = DebuggingRecorder::new(); let snapshotter = recorder.snapshotter(); @@ -293,7 +350,13 @@ mod tests { labels!("metastore_kind" => "read_replica", "test_label" => "test"), ); - primary_layer.record_request("hello", "transient", tonic::Code::Unavailable); + let mut retry_policy = primary_layer.retry_policy(RetryParams::for_test()); + let mut request = HelloRequest; + for num_attempts in 1..=3 { + let mut retryable_error: Result<(), RetryableError> = Err(RetryableError); + let retry = retry_policy.retry(&mut request, &mut retryable_error); + assert_eq!(retry.is_some(), num_attempts < 3); + } let mut hello_service = primary_layer.clone().layer(tower::service_fn( |request: HelloRequest| async move { Ok::<_, tonic::Status>(request) }, @@ -348,7 +411,7 @@ mod tests { ); assert_eq!( counter_value("hello", "transient", "unavailable", "primary"), - Some(&DebugValue::Counter(1)) + Some(&DebugValue::Counter(2)) ); assert_eq!( counter_value("goodbye", "success", "ok", "primary"), diff --git a/quickwit/quickwit-common/src/tower/mod.rs b/quickwit/quickwit-common/src/tower/mod.rs index d243d556e70..67748cf766b 100644 --- a/quickwit/quickwit-common/src/tower/mod.rs +++ b/quickwit/quickwit-common/src/tower/mod.rs @@ -44,13 +44,13 @@ 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}; pub use rate_estimator::{RateEstimator, SmaRateEstimator}; pub use rate_limit::{RateLimit, RateLimitLayer}; -pub use retry::{RetryCallbacks, RetryLayer, RetryPolicy}; +pub use retry::{RetryLayer, RetryPolicy}; pub use timeout::{Timeout, TimeoutExceeded, TimeoutLayer}; pub use transport::{BalanceChannel, warmup_channel}; diff --git a/quickwit/quickwit-common/src/tower/retry.rs b/quickwit/quickwit-common/src/tower/retry.rs index 7ed2435fb05..4594aaa1cda 100644 --- a/quickwit/quickwit-common/src/tower/retry.rs +++ b/quickwit/quickwit-common/src/tower/retry.rs @@ -25,8 +25,7 @@ use crate::retry::{RetryParams, Retryable}; /// Retry layer copy/pasted from `tower::retry::RetryLayer` /// but which implements `Clone`. impl Layer for RetryLayer

-where - P: Clone, +where P: Clone { type Service = Retry; @@ -48,40 +47,10 @@ impl

RetryLayer

{ } } -/// Callbacks invoked after the retry policy classifies an attempt's result. -pub trait RetryCallbacks: Clone { - /// Called when an attempt failed and another attempt will be made. - fn on_retry(&mut self, _request: &R, _error: &E, _num_attempts: usize) {} - - /// Called when the logical request completed with an error. - fn on_error(&mut self, _request: &R, _error: &E, _num_attempts: usize) {} - - /// Called when the logical request completed successfully. - fn on_success(&mut self, _request: &R, _response: &T, _num_attempts: usize) {} -} - -/// Default callbacks that do nothing. -#[derive(Clone, Copy, Debug, Default)] -pub struct NoopRetryCallbacks; - -impl RetryCallbacks for NoopRetryCallbacks {} - -#[derive(Clone, Debug)] -pub struct RetryPolicy { +#[derive(Clone, Copy, Debug)] +pub struct RetryPolicy { num_attempts: usize, retry_params: RetryParams, - callbacks: C, -} - -impl RetryPolicy { - /// Replaces the callbacks invoked when an attempt is classified. - pub fn with_callbacks(self, callbacks: D) -> RetryPolicy { - RetryPolicy { - num_attempts: self.num_attempts, - retry_params: self.retry_params, - callbacks, - } - } } impl From for RetryPolicy { @@ -89,34 +58,26 @@ impl From for RetryPolicy { Self { num_attempts: 0, retry_params, - callbacks: NoopRetryCallbacks, } } } -impl Policy for RetryPolicy +impl Policy for RetryPolicy where R: Clone, E: fmt::Debug + Retryable, - C: RetryCallbacks, { type Future = Sleep; - fn retry(&mut self, request: &mut R, result: &mut Result) -> Option { + fn retry(&mut self, _request: &mut R, result: &mut Result) -> Option { match result { - Ok(response) => { - self.callbacks - .on_success(request, response, self.num_attempts + 1); - None - } + Ok(_) => None, Err(error) => { self.num_attempts += 1; if !error.is_retryable() || self.num_attempts >= self.retry_params.max_attempts { - self.callbacks.on_error(request, error, self.num_attempts); None } else { - self.callbacks.on_retry(request, error, self.num_attempts); let delay = self.retry_params.compute_delay(self.num_attempts); debug!( num_attempts=%self.num_attempts, @@ -194,71 +155,6 @@ mod tests { } } - #[derive(Clone, Debug, Default)] - struct TestRetryCallbacks { - events: Arc>>, - } - - impl RetryCallbacks> for TestRetryCallbacks { - fn on_retry(&mut self, _request: &HelloRequest, _error: &Retry<()>, num_attempts: usize) { - self.events.lock().unwrap().push(("retry", num_attempts)); - } - - fn on_error(&mut self, _request: &HelloRequest, _error: &Retry<()>, num_attempts: usize) { - self.events.lock().unwrap().push(("error", num_attempts)); - } - - fn on_success(&mut self, _request: &HelloRequest, _response: &(), num_attempts: usize) { - self.events.lock().unwrap().push(("success", num_attempts)); - } - } - - #[tokio::test] - async fn test_retry_policy_callbacks() { - let callbacks = TestRetryCallbacks::default(); - let mut retry_policy = - RetryPolicy::from(RetryParams::for_test()).with_callbacks(callbacks.clone()); - let mut request = HelloRequest::default(); - - let mut transient_error = Err(Retry::Transient(())); - assert!( - retry_policy - .retry(&mut request, &mut transient_error) - .is_some() - ); - let mut success = Ok(()); - assert!(retry_policy.retry(&mut request, &mut success).is_none()); - assert_eq!( - callbacks.events.lock().unwrap().as_slice(), - &[("retry", 1), ("success", 2)] - ); - - let callbacks = TestRetryCallbacks::default(); - let mut retry_policy = - RetryPolicy::from(RetryParams::for_test()).with_callbacks(callbacks.clone()); - for num_attempts in 1..=3 { - let mut transient_error = Err(Retry::Transient(())); - let retry = retry_policy.retry(&mut request, &mut transient_error); - assert_eq!(retry.is_some(), num_attempts < 3); - } - assert_eq!( - callbacks.events.lock().unwrap().as_slice(), - &[("retry", 1), ("retry", 2), ("error", 3)] - ); - - let callbacks = TestRetryCallbacks::default(); - let mut retry_policy = - RetryPolicy::from(RetryParams::for_test()).with_callbacks(callbacks.clone()); - let mut permanent_error = Err(Retry::Permanent(())); - assert!(retry_policy - .retry(&mut request, &mut permanent_error) - .is_none()); - assert_eq!( - callbacks.events.lock().unwrap().as_slice(), - &[("error", 1)] - ); - } - #[tokio::test] async fn test_retry_policy() { let retry_policy = RetryPolicy::from(RetryParams::for_test()); diff --git a/quickwit/quickwit-serve/src/metastore.rs b/quickwit/quickwit-serve/src/metastore.rs index 9999a90593b..1ce7ba64c91 100644 --- a/quickwit/quickwit-serve/src/metastore.rs +++ b/quickwit/quickwit-serve/src/metastore.rs @@ -21,8 +21,7 @@ use quickwit_cluster::Cluster; use quickwit_common::pubsub::EventBroker; use quickwit_common::retry::RetryParams; use quickwit_common::tower::{ - EventListenerLayer, GrpcMetricsLayer, GrpcStatusCode, LoadShedLayer, RetryCallbacks, - RetryLayer, RetryPolicy, RpcName, TimeoutLayer, + EventListenerLayer, GrpcMetricsLayer, LoadShedLayer, RetryLayer, TimeoutLayer, }; use quickwit_common::uri::Uri; use quickwit_config::service::QuickwitService; @@ -72,22 +71,6 @@ static READ_REPLICA_METASTORE_GRPC_SERVER_METRICS_LAYER: LazyLock RetryCallbacks for MetastoreRetryCallbacks -where - R: RpcName, - E: GrpcStatusCode, -{ - fn on_retry(&mut self, _request: &R, error: &E, _num_attempts: usize) { - self.metrics_layer - .record_request(R::rpc_name(), "transient", error.grpc_status_code()); - } -} - fn get_metastore_client_max_concurrency() -> usize { quickwit_common::get_from_env( METASTORE_CLIENT_MAX_CONCURRENCY_ENV_KEY, @@ -193,15 +176,10 @@ impl LocalMetastoreServer { } _ => unreachable!("unexpected metastore service `{service}`"), }; - let retry_callbacks = MetastoreRetryCallbacks { - metrics_layer: metrics_layer.clone(), - }; - let retry_policy = - RetryPolicy::from(RetryParams::standard()).with_callbacks(retry_callbacks); + let retry_policy = metrics_layer.retry_policy(RetryParams::standard()); Ok(MetastoreServiceClient::tower() - // Keep the request metric outside retries so `success` and `error` represent the - // final result seen by the caller. RetryPolicy records each retryable failed attempt - // with `status="transient"`. + // Keep the metrics layer outside retries: the retry policy records each retryable + // failed attempt as `transient`, while the layer records the final outcome. .stack_layer(metrics_layer) .stack_layer(RetryLayer::new(retry_policy)) .stack_layer(TimeoutLayer::new(GRPC_METASTORE_SERVICE_TIMEOUT)) From 050b0066d345a504b7583f1a7c6208cf0aed3621 Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Mon, 17 Aug 2026 14:52:05 +0200 Subject: [PATCH 04/12] Revert "Simplify gRPC retry metrics" This reverts commit ff93535c5ab4a7001eb5f2b33decd4e3acfa70bb. --- quickwit/quickwit-common/src/tower/metrics.rs | 79 ++---------- quickwit/quickwit-common/src/tower/mod.rs | 4 +- quickwit/quickwit-common/src/tower/retry.rs | 116 +++++++++++++++++- quickwit/quickwit-serve/src/metastore.rs | 30 ++++- 4 files changed, 146 insertions(+), 83 deletions(-) diff --git a/quickwit/quickwit-common/src/tower/metrics.rs b/quickwit/quickwit-common/src/tower/metrics.rs index 3d920f33379..8ef18cbd421 100644 --- a/quickwit/quickwit-common/src/tower/metrics.rs +++ b/quickwit/quickwit-common/src/tower/metrics.rs @@ -12,7 +12,6 @@ // See the License for the specific language governing permissions and // limitations under the License. -use std::fmt; use std::pin::Pin; use std::task::{Context, Poll}; use std::time::Instant; @@ -23,12 +22,9 @@ 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 super::retry::RetryPolicy; use crate::metrics::exponential_buckets; -use crate::retry::{RetryParams, Retryable}; pub trait RpcName { fn rpc_name() -> &'static str; @@ -51,7 +47,7 @@ impl GrpcStatusCode for std::convert::Infallible { } } -fn grpc_code_label(code: tonic::Code) -> &'static str { +pub(crate) fn grpc_code_label(code: tonic::Code) -> &'static str { match code { tonic::Code::Ok => "ok", tonic::Code::Cancelled => "cancelled", @@ -145,37 +141,6 @@ pub struct GrpcMetricsLayer { request_duration_seconds: Histogram, } -/// Retry policy that records each failed attempt which will be retried as transient. -#[derive(Clone)] -pub struct GrpcRetryPolicy { - inner: RetryPolicy, - metrics_layer: GrpcMetricsLayer, -} - -impl Policy for GrpcRetryPolicy -where - R: Clone + RpcName, - E: fmt::Debug + Retryable + GrpcStatusCode, -{ - type Future = >::Future; - - fn retry(&mut self, request: &mut R, result: &mut Result) -> Option { - let retry = self.inner.retry(request, result); - if retry.is_some() { - let Err(error) = result else { - unreachable!("a successful request should not be retried"); - }; - self.metrics_layer - .record_request(R::rpc_name(), "transient", error.grpc_status_code()); - } - retry - } - - fn clone_request(&mut self, request: &R) -> Option { - >::clone_request(&mut self.inner, request) - } -} - impl GrpcMetricsLayer { pub fn new(subsystem: &'static str, kind: &'static str) -> Self { let labels = Self::default_labels(subsystem, kind); @@ -199,15 +164,8 @@ impl GrpcMetricsLayer { } } - /// Wraps the retry policy to record retryable failed attempts as transient requests. - pub fn retry_policy(&self, retry_params: RetryParams) -> GrpcRetryPolicy { - GrpcRetryPolicy { - inner: RetryPolicy::from(retry_params), - metrics_layer: self.clone(), - } - } - - fn record_request(&self, rpc_name: &'static str, status: &'static str, code: tonic::Code) { + /// Records a request outcome with this layer's static labels. + pub fn record_request(&self, rpc_name: &'static str, status: &'static str, code: tonic::Code) { counter!( parent: self.requests_total, labels: [labels!( @@ -300,7 +258,7 @@ mod tests { use super::*; - #[derive(Clone, Debug)] + #[derive(Debug)] struct HelloRequest; impl RpcName for HelloRequest { @@ -311,29 +269,14 @@ mod tests { struct GoodbyeRequest; - #[derive(Debug)] - struct RetryableError; - - impl Retryable for RetryableError { - fn is_retryable(&self) -> bool { - true - } - } - - impl GrpcStatusCode for RetryableError { - fn grpc_status_code(&self) -> tonic::Code { - tonic::Code::Unavailable - } - } - impl RpcName for GoodbyeRequest { fn rpc_name() -> &'static str { "goodbye" } } - #[tokio::test] - async fn test_grpc_metrics() { + #[test] + fn test_grpc_metrics() { let recorder = DebuggingRecorder::new(); let snapshotter = recorder.snapshotter(); @@ -350,13 +293,7 @@ mod tests { labels!("metastore_kind" => "read_replica", "test_label" => "test"), ); - let mut retry_policy = primary_layer.retry_policy(RetryParams::for_test()); - let mut request = HelloRequest; - for num_attempts in 1..=3 { - let mut retryable_error: Result<(), RetryableError> = Err(RetryableError); - let retry = retry_policy.retry(&mut request, &mut retryable_error); - assert_eq!(retry.is_some(), num_attempts < 3); - } + primary_layer.record_request("hello", "transient", tonic::Code::Unavailable); let mut hello_service = primary_layer.clone().layer(tower::service_fn( |request: HelloRequest| async move { Ok::<_, tonic::Status>(request) }, @@ -411,7 +348,7 @@ mod tests { ); assert_eq!( counter_value("hello", "transient", "unavailable", "primary"), - Some(&DebugValue::Counter(2)) + Some(&DebugValue::Counter(1)) ); assert_eq!( counter_value("goodbye", "success", "ok", "primary"), diff --git a/quickwit/quickwit-common/src/tower/mod.rs b/quickwit/quickwit-common/src/tower/mod.rs index 67748cf766b..d243d556e70 100644 --- a/quickwit/quickwit-common/src/tower/mod.rs +++ b/quickwit/quickwit-common/src/tower/mod.rs @@ -44,13 +44,13 @@ 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, GrpcRetryPolicy, GrpcStatusCode, RpcName}; +pub use metrics::{GrpcMetrics, GrpcMetricsLayer, GrpcStatusCode, RpcName}; pub use one_task_per_call_layer::{OneTaskPerCallLayer, TaskCancelled}; pub use pool::Pool; pub use rate::{ConstantRate, Rate}; pub use rate_estimator::{RateEstimator, SmaRateEstimator}; pub use rate_limit::{RateLimit, RateLimitLayer}; -pub use retry::{RetryLayer, RetryPolicy}; +pub use retry::{RetryCallbacks, RetryLayer, RetryPolicy}; pub use timeout::{Timeout, TimeoutExceeded, TimeoutLayer}; pub use transport::{BalanceChannel, warmup_channel}; diff --git a/quickwit/quickwit-common/src/tower/retry.rs b/quickwit/quickwit-common/src/tower/retry.rs index 4594aaa1cda..7ed2435fb05 100644 --- a/quickwit/quickwit-common/src/tower/retry.rs +++ b/quickwit/quickwit-common/src/tower/retry.rs @@ -25,7 +25,8 @@ use crate::retry::{RetryParams, Retryable}; /// Retry layer copy/pasted from `tower::retry::RetryLayer` /// but which implements `Clone`. impl Layer for RetryLayer

-where P: Clone +where + P: Clone, { type Service = Retry; @@ -47,10 +48,40 @@ impl

RetryLayer

{ } } -#[derive(Clone, Copy, Debug)] -pub struct RetryPolicy { +/// Callbacks invoked after the retry policy classifies an attempt's result. +pub trait RetryCallbacks: Clone { + /// Called when an attempt failed and another attempt will be made. + fn on_retry(&mut self, _request: &R, _error: &E, _num_attempts: usize) {} + + /// Called when the logical request completed with an error. + fn on_error(&mut self, _request: &R, _error: &E, _num_attempts: usize) {} + + /// Called when the logical request completed successfully. + fn on_success(&mut self, _request: &R, _response: &T, _num_attempts: usize) {} +} + +/// Default callbacks that do nothing. +#[derive(Clone, Copy, Debug, Default)] +pub struct NoopRetryCallbacks; + +impl RetryCallbacks for NoopRetryCallbacks {} + +#[derive(Clone, Debug)] +pub struct RetryPolicy { num_attempts: usize, retry_params: RetryParams, + callbacks: C, +} + +impl RetryPolicy { + /// Replaces the callbacks invoked when an attempt is classified. + pub fn with_callbacks(self, callbacks: D) -> RetryPolicy { + RetryPolicy { + num_attempts: self.num_attempts, + retry_params: self.retry_params, + callbacks, + } + } } impl From for RetryPolicy { @@ -58,26 +89,34 @@ impl From for RetryPolicy { Self { num_attempts: 0, retry_params, + callbacks: NoopRetryCallbacks, } } } -impl Policy for RetryPolicy +impl Policy for RetryPolicy where R: Clone, E: fmt::Debug + Retryable, + C: RetryCallbacks, { type Future = Sleep; - fn retry(&mut self, _request: &mut R, result: &mut Result) -> Option { + fn retry(&mut self, request: &mut R, result: &mut Result) -> Option { match result { - Ok(_) => None, + Ok(response) => { + self.callbacks + .on_success(request, response, self.num_attempts + 1); + None + } Err(error) => { self.num_attempts += 1; if !error.is_retryable() || self.num_attempts >= self.retry_params.max_attempts { + self.callbacks.on_error(request, error, self.num_attempts); None } else { + self.callbacks.on_retry(request, error, self.num_attempts); let delay = self.retry_params.compute_delay(self.num_attempts); debug!( num_attempts=%self.num_attempts, @@ -155,6 +194,71 @@ mod tests { } } + #[derive(Clone, Debug, Default)] + struct TestRetryCallbacks { + events: Arc>>, + } + + impl RetryCallbacks> for TestRetryCallbacks { + fn on_retry(&mut self, _request: &HelloRequest, _error: &Retry<()>, num_attempts: usize) { + self.events.lock().unwrap().push(("retry", num_attempts)); + } + + fn on_error(&mut self, _request: &HelloRequest, _error: &Retry<()>, num_attempts: usize) { + self.events.lock().unwrap().push(("error", num_attempts)); + } + + fn on_success(&mut self, _request: &HelloRequest, _response: &(), num_attempts: usize) { + self.events.lock().unwrap().push(("success", num_attempts)); + } + } + + #[tokio::test] + async fn test_retry_policy_callbacks() { + let callbacks = TestRetryCallbacks::default(); + let mut retry_policy = + RetryPolicy::from(RetryParams::for_test()).with_callbacks(callbacks.clone()); + let mut request = HelloRequest::default(); + + let mut transient_error = Err(Retry::Transient(())); + assert!( + retry_policy + .retry(&mut request, &mut transient_error) + .is_some() + ); + let mut success = Ok(()); + assert!(retry_policy.retry(&mut request, &mut success).is_none()); + assert_eq!( + callbacks.events.lock().unwrap().as_slice(), + &[("retry", 1), ("success", 2)] + ); + + let callbacks = TestRetryCallbacks::default(); + let mut retry_policy = + RetryPolicy::from(RetryParams::for_test()).with_callbacks(callbacks.clone()); + for num_attempts in 1..=3 { + let mut transient_error = Err(Retry::Transient(())); + let retry = retry_policy.retry(&mut request, &mut transient_error); + assert_eq!(retry.is_some(), num_attempts < 3); + } + assert_eq!( + callbacks.events.lock().unwrap().as_slice(), + &[("retry", 1), ("retry", 2), ("error", 3)] + ); + + let callbacks = TestRetryCallbacks::default(); + let mut retry_policy = + RetryPolicy::from(RetryParams::for_test()).with_callbacks(callbacks.clone()); + let mut permanent_error = Err(Retry::Permanent(())); + assert!(retry_policy + .retry(&mut request, &mut permanent_error) + .is_none()); + assert_eq!( + callbacks.events.lock().unwrap().as_slice(), + &[("error", 1)] + ); + } + #[tokio::test] async fn test_retry_policy() { let retry_policy = RetryPolicy::from(RetryParams::for_test()); diff --git a/quickwit/quickwit-serve/src/metastore.rs b/quickwit/quickwit-serve/src/metastore.rs index 1ce7ba64c91..9999a90593b 100644 --- a/quickwit/quickwit-serve/src/metastore.rs +++ b/quickwit/quickwit-serve/src/metastore.rs @@ -21,7 +21,8 @@ use quickwit_cluster::Cluster; use quickwit_common::pubsub::EventBroker; use quickwit_common::retry::RetryParams; use quickwit_common::tower::{ - EventListenerLayer, GrpcMetricsLayer, LoadShedLayer, RetryLayer, TimeoutLayer, + EventListenerLayer, GrpcMetricsLayer, GrpcStatusCode, LoadShedLayer, RetryCallbacks, + RetryLayer, RetryPolicy, RpcName, TimeoutLayer, }; use quickwit_common::uri::Uri; use quickwit_config::service::QuickwitService; @@ -71,6 +72,22 @@ static READ_REPLICA_METASTORE_GRPC_SERVER_METRICS_LAYER: LazyLock RetryCallbacks for MetastoreRetryCallbacks +where + R: RpcName, + E: GrpcStatusCode, +{ + fn on_retry(&mut self, _request: &R, error: &E, _num_attempts: usize) { + self.metrics_layer + .record_request(R::rpc_name(), "transient", error.grpc_status_code()); + } +} + fn get_metastore_client_max_concurrency() -> usize { quickwit_common::get_from_env( METASTORE_CLIENT_MAX_CONCURRENCY_ENV_KEY, @@ -176,10 +193,15 @@ impl LocalMetastoreServer { } _ => unreachable!("unexpected metastore service `{service}`"), }; - let retry_policy = metrics_layer.retry_policy(RetryParams::standard()); + let retry_callbacks = MetastoreRetryCallbacks { + metrics_layer: metrics_layer.clone(), + }; + let retry_policy = + RetryPolicy::from(RetryParams::standard()).with_callbacks(retry_callbacks); Ok(MetastoreServiceClient::tower() - // Keep the metrics layer outside retries: the retry policy records each retryable - // failed attempt as `transient`, while the layer records the final outcome. + // Keep the request metric outside retries so `success` and `error` represent the + // final result seen by the caller. RetryPolicy records each retryable failed attempt + // with `status="transient"`. .stack_layer(metrics_layer) .stack_layer(RetryLayer::new(retry_policy)) .stack_layer(TimeoutLayer::new(GRPC_METASTORE_SERVICE_TIMEOUT)) From 0174d9ecf6713051c083e5770a7a0f4b09c0910d Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Mon, 17 Aug 2026 14:52:05 +0200 Subject: [PATCH 05/12] Revert "Decouple retry metrics with callbacks" This reverts commit b2f1fc9942fe9fcb6fa4424065094e1ae5480625. --- quickwit/quickwit-common/src/tower/metrics.rs | 23 +-- quickwit/quickwit-common/src/tower/mod.rs | 2 +- quickwit/quickwit-common/src/tower/retry.rs | 177 +++++++----------- quickwit/quickwit-serve/src/metastore.rs | 26 +-- 4 files changed, 82 insertions(+), 146 deletions(-) diff --git a/quickwit/quickwit-common/src/tower/metrics.rs b/quickwit/quickwit-common/src/tower/metrics.rs index 8ef18cbd421..71db523e9ed 100644 --- a/quickwit/quickwit-common/src/tower/metrics.rs +++ b/quickwit/quickwit-common/src/tower/metrics.rs @@ -164,17 +164,12 @@ impl GrpcMetricsLayer { } } - /// Records a request outcome with this layer's static labels. - pub 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(); + /// Returns the request counter with this layer's static labels. + /// + /// A retry policy uses this to record a retryable failed attempt with + /// `status="transient"` after it decides to retry it. + pub fn requests_total_counter(&self) -> Counter { + self.requests_total.clone() } fn default_labels(subsystem: &'static str, kind: &'static str) -> Labels<3> { @@ -293,8 +288,6 @@ mod tests { labels!("metastore_kind" => "read_replica", "test_label" => "test"), ); - primary_layer.record_request("hello", "transient", tonic::Code::Unavailable); - let mut hello_service = primary_layer.clone().layer(tower::service_fn( |request: HelloRequest| async move { Ok::<_, tonic::Status>(request) }, )); @@ -346,10 +339,6 @@ mod tests { counter_value("hello", "success", "ok", "primary"), Some(&DebugValue::Counter(1)) ); - assert_eq!( - counter_value("hello", "transient", "unavailable", "primary"), - Some(&DebugValue::Counter(1)) - ); assert_eq!( counter_value("goodbye", "success", "ok", "primary"), Some(&DebugValue::Counter(1)) diff --git a/quickwit/quickwit-common/src/tower/mod.rs b/quickwit/quickwit-common/src/tower/mod.rs index d243d556e70..c763ce22bf8 100644 --- a/quickwit/quickwit-common/src/tower/mod.rs +++ b/quickwit/quickwit-common/src/tower/mod.rs @@ -50,7 +50,7 @@ pub use pool::Pool; pub use rate::{ConstantRate, Rate}; pub use rate_estimator::{RateEstimator, SmaRateEstimator}; pub use rate_limit::{RateLimit, RateLimitLayer}; -pub use retry::{RetryCallbacks, RetryLayer, RetryPolicy}; +pub use retry::{RetryLayer, RetryPolicy}; pub use timeout::{Timeout, TimeoutExceeded, TimeoutLayer}; pub use transport::{BalanceChannel, warmup_channel}; diff --git a/quickwit/quickwit-common/src/tower/retry.rs b/quickwit/quickwit-common/src/tower/retry.rs index 7ed2435fb05..a07bd2d0701 100644 --- a/quickwit/quickwit-common/src/tower/retry.rs +++ b/quickwit/quickwit-common/src/tower/retry.rs @@ -15,18 +15,19 @@ use std::any::type_name; use std::fmt; +use quickwit_metrics::{Counter, counter, labels}; use tokio::time::Sleep; use tower::Layer; use tower::retry::{Policy, Retry}; use tracing::debug; +use super::metrics::{GrpcStatusCode, RpcName, grpc_code_label}; use crate::retry::{RetryParams, Retryable}; /// Retry layer copy/pasted from `tower::retry::RetryLayer` /// but which implements `Clone`. impl Layer for RetryLayer

-where - P: Clone, +where P: Clone { type Service = Retry; @@ -48,39 +49,18 @@ impl

RetryLayer

{ } } -/// Callbacks invoked after the retry policy classifies an attempt's result. -pub trait RetryCallbacks: Clone { - /// Called when an attempt failed and another attempt will be made. - fn on_retry(&mut self, _request: &R, _error: &E, _num_attempts: usize) {} - - /// Called when the logical request completed with an error. - fn on_error(&mut self, _request: &R, _error: &E, _num_attempts: usize) {} - - /// Called when the logical request completed successfully. - fn on_success(&mut self, _request: &R, _response: &T, _num_attempts: usize) {} -} - -/// Default callbacks that do nothing. -#[derive(Clone, Copy, Debug, Default)] -pub struct NoopRetryCallbacks; - -impl RetryCallbacks for NoopRetryCallbacks {} - #[derive(Clone, Debug)] -pub struct RetryPolicy { +pub struct RetryPolicy { num_attempts: usize, retry_params: RetryParams, - callbacks: C, + retry_metrics_counter_opt: Option, } -impl RetryPolicy { - /// Replaces the callbacks invoked when an attempt is classified. - pub fn with_callbacks(self, callbacks: D) -> RetryPolicy { - RetryPolicy { - num_attempts: self.num_attempts, - retry_params: self.retry_params, - callbacks, - } +impl RetryPolicy { + /// Records each failed attempt that this policy will retry as a transient request. + pub fn with_retry_metrics(mut self, retry_metrics_counter: Counter) -> Self { + self.retry_metrics_counter_opt = Some(retry_metrics_counter); + self } } @@ -89,34 +69,38 @@ impl From for RetryPolicy { Self { num_attempts: 0, retry_params, - callbacks: NoopRetryCallbacks, + retry_metrics_counter_opt: None, } } } -impl Policy for RetryPolicy +impl Policy for RetryPolicy where - R: Clone, - E: fmt::Debug + Retryable, - C: RetryCallbacks, + R: Clone + RpcName, + E: fmt::Debug + Retryable + GrpcStatusCode, { type Future = Sleep; - fn retry(&mut self, request: &mut R, result: &mut Result) -> Option { + fn retry(&mut self, _request: &mut R, result: &mut Result) -> Option { match result { - Ok(response) => { - self.callbacks - .on_success(request, response, self.num_attempts + 1); - None - } + Ok(_) => None, Err(error) => { self.num_attempts += 1; if !error.is_retryable() || self.num_attempts >= self.retry_params.max_attempts { - self.callbacks.on_error(request, error, self.num_attempts); None } else { - self.callbacks.on_retry(request, error, self.num_attempts); + if let Some(retry_metrics_counter) = &self.retry_metrics_counter_opt { + counter!( + parent: retry_metrics_counter, + labels: [labels!( + "rpc" => R::rpc_name(), + "status" => "transient", + "code" => grpc_code_label(error.grpc_status_code()), + )], + ) + .inc(); + } let delay = self.retry_params.compute_delay(self.num_attempts); debug!( num_attempts=%self.num_attempts, @@ -143,6 +127,8 @@ mod tests { use std::task::{Context, Poll}; use futures::future::{Ready, ready}; + use metrics::with_local_recorder; + use metrics_util::debugging::{DebugValue, DebuggingRecorder}; use tower::{Layer, Service, ServiceExt}; use super::*; @@ -162,6 +148,12 @@ mod tests { } } + impl GrpcStatusCode for Retry { + fn grpc_status_code(&self) -> tonic::Code { + tonic::Code::Unavailable + } + } + #[derive(Debug, Clone, Default)] struct HelloService; @@ -173,6 +165,12 @@ mod tests { results: HelloResults, } + impl RpcName for HelloRequest { + fn rpc_name() -> &'static str { + "hello" + } + } + impl Service for HelloService { type Response = (); type Error = Retry<()>; @@ -194,69 +192,38 @@ mod tests { } } - #[derive(Clone, Debug, Default)] - struct TestRetryCallbacks { - events: Arc>>, - } - - impl RetryCallbacks> for TestRetryCallbacks { - fn on_retry(&mut self, _request: &HelloRequest, _error: &Retry<()>, num_attempts: usize) { - self.events.lock().unwrap().push(("retry", num_attempts)); - } - - fn on_error(&mut self, _request: &HelloRequest, _error: &Retry<()>, num_attempts: usize) { - self.events.lock().unwrap().push(("error", num_attempts)); - } - - fn on_success(&mut self, _request: &HelloRequest, _response: &(), num_attempts: usize) { - self.events.lock().unwrap().push(("success", num_attempts)); - } - } - #[tokio::test] - async fn test_retry_policy_callbacks() { - let callbacks = TestRetryCallbacks::default(); - let mut retry_policy = - RetryPolicy::from(RetryParams::for_test()).with_callbacks(callbacks.clone()); - let mut request = HelloRequest::default(); - - let mut transient_error = Err(Retry::Transient(())); - assert!( - retry_policy - .retry(&mut request, &mut transient_error) - .is_some() - ); - let mut success = Ok(()); - assert!(retry_policy.retry(&mut request, &mut success).is_none()); - assert_eq!( - callbacks.events.lock().unwrap().as_slice(), - &[("retry", 1), ("success", 2)] - ); - - let callbacks = TestRetryCallbacks::default(); - let mut retry_policy = - RetryPolicy::from(RetryParams::for_test()).with_callbacks(callbacks.clone()); - for num_attempts in 1..=3 { - let mut transient_error = Err(Retry::Transient(())); - let retry = retry_policy.retry(&mut request, &mut transient_error); - assert_eq!(retry.is_some(), num_attempts < 3); - } - assert_eq!( - callbacks.events.lock().unwrap().as_slice(), - &[("retry", 1), ("retry", 2), ("error", 3)] - ); - - let callbacks = TestRetryCallbacks::default(); - let mut retry_policy = - RetryPolicy::from(RetryParams::for_test()).with_callbacks(callbacks.clone()); - let mut permanent_error = Err(Retry::Permanent(())); - assert!(retry_policy - .retry(&mut request, &mut permanent_error) - .is_none()); - assert_eq!( - callbacks.events.lock().unwrap().as_slice(), - &[("error", 1)] - ); + async fn test_retry_policy_records_retryable_failures_as_transient() { + let recorder = DebuggingRecorder::new(); + let snapshotter = recorder.snapshotter(); + + with_local_recorder(&recorder, || { + let retry_metrics_counter = counter!( + name: "requests_total", + description: "test request count", + subsystem: "grpc", + ); + let mut retry_policy = RetryPolicy::from(RetryParams::for_test()) + .with_retry_metrics(retry_metrics_counter); + let mut request = HelloRequest::default(); + let mut result: Result<(), Retry<()>> = Err(Retry::Transient(())); + + assert!(retry_policy.retry(&mut request, &mut result).is_some()); + }); + + let snapshot = snapshotter.snapshot().into_vec(); + assert!(snapshot.iter().any(|(composite_key, _, _, value)| { + let (_, key) = composite_key.clone().into_parts(); + let labels = key + .labels() + .map(|label| (label.key(), label.value())) + .collect::>(); + key.name() == "quickwit_grpc_requests_total" + && labels.contains(&("rpc", "hello")) + && labels.contains(&("status", "transient")) + && labels.contains(&("code", "unavailable")) + && value == &DebugValue::Counter(1) + })); } #[tokio::test] diff --git a/quickwit/quickwit-serve/src/metastore.rs b/quickwit/quickwit-serve/src/metastore.rs index 9999a90593b..5b9a693e1c6 100644 --- a/quickwit/quickwit-serve/src/metastore.rs +++ b/quickwit/quickwit-serve/src/metastore.rs @@ -21,8 +21,7 @@ use quickwit_cluster::Cluster; use quickwit_common::pubsub::EventBroker; use quickwit_common::retry::RetryParams; use quickwit_common::tower::{ - EventListenerLayer, GrpcMetricsLayer, GrpcStatusCode, LoadShedLayer, RetryCallbacks, - RetryLayer, RetryPolicy, RpcName, TimeoutLayer, + EventListenerLayer, GrpcMetricsLayer, LoadShedLayer, RetryLayer, RetryPolicy, TimeoutLayer, }; use quickwit_common::uri::Uri; use quickwit_config::service::QuickwitService; @@ -72,22 +71,6 @@ static READ_REPLICA_METASTORE_GRPC_SERVER_METRICS_LAYER: LazyLock RetryCallbacks for MetastoreRetryCallbacks -where - R: RpcName, - E: GrpcStatusCode, -{ - fn on_retry(&mut self, _request: &R, error: &E, _num_attempts: usize) { - self.metrics_layer - .record_request(R::rpc_name(), "transient", error.grpc_status_code()); - } -} - fn get_metastore_client_max_concurrency() -> usize { quickwit_common::get_from_env( METASTORE_CLIENT_MAX_CONCURRENCY_ENV_KEY, @@ -193,11 +176,8 @@ impl LocalMetastoreServer { } _ => unreachable!("unexpected metastore service `{service}`"), }; - let retry_callbacks = MetastoreRetryCallbacks { - metrics_layer: metrics_layer.clone(), - }; - let retry_policy = - RetryPolicy::from(RetryParams::standard()).with_callbacks(retry_callbacks); + let retry_policy = RetryPolicy::from(RetryParams::standard()) + .with_retry_metrics(metrics_layer.requests_total_counter()); Ok(MetastoreServiceClient::tower() // Keep the request metric outside retries so `success` and `error` represent the // final result seen by the caller. RetryPolicy records each retryable failed attempt From cdbfdcfbddbfb2f157916412c4f88cd493f3f38c Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Mon, 17 Aug 2026 14:52:06 +0200 Subject: [PATCH 06/12] Revert "Add metastore retry outcome metrics" This reverts commit cf343f7815a3fd291ae6e7fb938a1292ade3bc89. --- docs/reference/metrics.md | 19 ++--- quickwit/quickwit-common/src/tower/metrics.rs | 10 +-- quickwit/quickwit-common/src/tower/retry.rs | 77 +------------------ quickwit/quickwit-serve/src/metastore.rs | 9 +-- 4 files changed, 11 insertions(+), 104 deletions(-) diff --git a/docs/reference/metrics.md b/docs/reference/metrics.md index d28984e8745..33a49854895 100644 --- a/docs/reference/metrics.md +++ b/docs/reference/metrics.md @@ -50,24 +50,15 @@ Quickwit exposes metrics for several cache components, including `fastfields`, ` ## Metastore Metrics -Metastore RPCs use the shared gRPC metrics: +All metastore methods are monitored by the 3 metrics: | Namespace | Metric Name | Description | Labels | Type | | --------- | ----------- | ----------- | ------ | ---- | -| `quickwit_grpc` | `requests_total` | Number of metastore RPC outcomes and retryable failed attempts | [`grpc_service`, `kind`, `metastore_kind`, `rpc`, `status`, `code`] | `counter` | +| `quickwit_metastore` | `requests_total` | Number of requests | [`operation`, `index`] | `counter` | +| `quickwit_metastore` | `request_errors_total` | Number of failed requests | [`operation`, `index`] | `counter` | +| `quickwit_metastore` | `request_duration_seconds` | Duration of requests | [`operation`, `index`, `error`] | `histogram` | -Filter this metric with `grpc_service="metastore"`. `kind` identifies the client or server, -and `metastore_kind` identifies the primary or read-replica metastore. The `status` label has the -following meaning for client metrics: - -- `success`: the logical RPC succeeded, including after any retries. -- `transient`: a retryable failed attempt for which the client will issue another attempt. -- `error`: the logical RPC failed and was returned to the caller (a non-retryable error or after - retry exhaustion). -- `cancelled`: the client cancelled the RPC before it completed. - -`rpc` contains the metastore operation name, such as `create_index`, `index_metadata`, -`delete_index`, `stage_splits`, `publish_splits`, `list_splits`, or `add_source`. +Examples of operation names: `create_index`, `index_metadata`, `delete_index`, `stage_splits`, `publish_splits`, `list_splits`, `add_source`, ... PostgreSQL-backed metastores also expose connection pool gauges: diff --git a/quickwit/quickwit-common/src/tower/metrics.rs b/quickwit/quickwit-common/src/tower/metrics.rs index 71db523e9ed..17a9f6acec3 100644 --- a/quickwit/quickwit-common/src/tower/metrics.rs +++ b/quickwit/quickwit-common/src/tower/metrics.rs @@ -47,7 +47,7 @@ impl GrpcStatusCode for std::convert::Infallible { } } -pub(crate) fn grpc_code_label(code: tonic::Code) -> &'static str { +fn grpc_code_label(code: tonic::Code) -> &'static str { match code { tonic::Code::Ok => "ok", tonic::Code::Cancelled => "cancelled", @@ -164,14 +164,6 @@ impl GrpcMetricsLayer { } } - /// Returns the request counter with this layer's static labels. - /// - /// A retry policy uses this to record a retryable failed attempt with - /// `status="transient"` after it decides to retry it. - pub fn requests_total_counter(&self) -> Counter { - self.requests_total.clone() - } - 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. diff --git a/quickwit/quickwit-common/src/tower/retry.rs b/quickwit/quickwit-common/src/tower/retry.rs index a07bd2d0701..4594aaa1cda 100644 --- a/quickwit/quickwit-common/src/tower/retry.rs +++ b/quickwit/quickwit-common/src/tower/retry.rs @@ -15,13 +15,11 @@ use std::any::type_name; use std::fmt; -use quickwit_metrics::{Counter, counter, labels}; use tokio::time::Sleep; use tower::Layer; use tower::retry::{Policy, Retry}; use tracing::debug; -use super::metrics::{GrpcStatusCode, RpcName, grpc_code_label}; use crate::retry::{RetryParams, Retryable}; /// Retry layer copy/pasted from `tower::retry::RetryLayer` @@ -49,19 +47,10 @@ impl

RetryLayer

{ } } -#[derive(Clone, Debug)] +#[derive(Clone, Copy, Debug)] pub struct RetryPolicy { num_attempts: usize, retry_params: RetryParams, - retry_metrics_counter_opt: Option, -} - -impl RetryPolicy { - /// Records each failed attempt that this policy will retry as a transient request. - pub fn with_retry_metrics(mut self, retry_metrics_counter: Counter) -> Self { - self.retry_metrics_counter_opt = Some(retry_metrics_counter); - self - } } impl From for RetryPolicy { @@ -69,15 +58,14 @@ impl From for RetryPolicy { Self { num_attempts: 0, retry_params, - retry_metrics_counter_opt: None, } } } impl Policy for RetryPolicy where - R: Clone + RpcName, - E: fmt::Debug + Retryable + GrpcStatusCode, + R: Clone, + E: fmt::Debug + Retryable, { type Future = Sleep; @@ -90,17 +78,6 @@ where if !error.is_retryable() || self.num_attempts >= self.retry_params.max_attempts { None } else { - if let Some(retry_metrics_counter) = &self.retry_metrics_counter_opt { - counter!( - parent: retry_metrics_counter, - labels: [labels!( - "rpc" => R::rpc_name(), - "status" => "transient", - "code" => grpc_code_label(error.grpc_status_code()), - )], - ) - .inc(); - } let delay = self.retry_params.compute_delay(self.num_attempts); debug!( num_attempts=%self.num_attempts, @@ -127,8 +104,6 @@ mod tests { use std::task::{Context, Poll}; use futures::future::{Ready, ready}; - use metrics::with_local_recorder; - use metrics_util::debugging::{DebugValue, DebuggingRecorder}; use tower::{Layer, Service, ServiceExt}; use super::*; @@ -148,12 +123,6 @@ mod tests { } } - impl GrpcStatusCode for Retry { - fn grpc_status_code(&self) -> tonic::Code { - tonic::Code::Unavailable - } - } - #[derive(Debug, Clone, Default)] struct HelloService; @@ -165,12 +134,6 @@ mod tests { results: HelloResults, } - impl RpcName for HelloRequest { - fn rpc_name() -> &'static str { - "hello" - } - } - impl Service for HelloService { type Response = (); type Error = Retry<()>; @@ -192,40 +155,6 @@ mod tests { } } - #[tokio::test] - async fn test_retry_policy_records_retryable_failures_as_transient() { - let recorder = DebuggingRecorder::new(); - let snapshotter = recorder.snapshotter(); - - with_local_recorder(&recorder, || { - let retry_metrics_counter = counter!( - name: "requests_total", - description: "test request count", - subsystem: "grpc", - ); - let mut retry_policy = RetryPolicy::from(RetryParams::for_test()) - .with_retry_metrics(retry_metrics_counter); - let mut request = HelloRequest::default(); - let mut result: Result<(), Retry<()>> = Err(Retry::Transient(())); - - assert!(retry_policy.retry(&mut request, &mut result).is_some()); - }); - - let snapshot = snapshotter.snapshot().into_vec(); - assert!(snapshot.iter().any(|(composite_key, _, _, value)| { - let (_, key) = composite_key.clone().into_parts(); - let labels = key - .labels() - .map(|label| (label.key(), label.value())) - .collect::>(); - key.name() == "quickwit_grpc_requests_total" - && labels.contains(&("rpc", "hello")) - && labels.contains(&("status", "transient")) - && labels.contains(&("code", "unavailable")) - && value == &DebugValue::Counter(1) - })); - } - #[tokio::test] async fn test_retry_policy() { let retry_policy = RetryPolicy::from(RetryParams::for_test()); diff --git a/quickwit/quickwit-serve/src/metastore.rs b/quickwit/quickwit-serve/src/metastore.rs index 5b9a693e1c6..a4ee434232b 100644 --- a/quickwit/quickwit-serve/src/metastore.rs +++ b/quickwit/quickwit-serve/src/metastore.rs @@ -176,15 +176,10 @@ impl LocalMetastoreServer { } _ => unreachable!("unexpected metastore service `{service}`"), }; - let retry_policy = RetryPolicy::from(RetryParams::standard()) - .with_retry_metrics(metrics_layer.requests_total_counter()); Ok(MetastoreServiceClient::tower() - // Keep the request metric outside retries so `success` and `error` represent the - // final result seen by the caller. RetryPolicy records each retryable failed attempt - // with `status="transient"`. - .stack_layer(metrics_layer) - .stack_layer(RetryLayer::new(retry_policy)) + .stack_layer(RetryLayer::new(RetryPolicy::from(RetryParams::standard()))) .stack_layer(TimeoutLayer::new(GRPC_METASTORE_SERVICE_TIMEOUT)) + .stack_layer(metrics_layer) .stack_layer(tower::limit::GlobalConcurrencyLimitLayer::new( get_metastore_client_max_concurrency(), )) From 1c193b34ea3dc0b5885da11894eeb26e338cbc8a Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Mon, 17 Aug 2026 15:56:10 +0200 Subject: [PATCH 07/12] Distinguish metastore retry attempts from errors Record scheduled retries separately so final error metrics reflect only client-visible failures. --- docs/reference/metrics.md | 10 + quickwit/quickwit-common/src/tower/metrics.rs | 173 +++++++++++++++++- quickwit/quickwit-common/src/tower/mod.rs | 2 +- quickwit/quickwit-serve/src/metastore.rs | 9 +- 4 files changed, 180 insertions(+), 14 deletions(-) diff --git a/docs/reference/metrics.md b/docs/reference/metrics.md index 33a49854895..293903fcb76 100644 --- a/docs/reference/metrics.md +++ b/docs/reference/metrics.md @@ -69,6 +69,16 @@ PostgreSQL-backed metastores also expose connection pool gauges: | `quickwit_metastore` | `acquire_connections` | Number of requests currently waiting to acquire a PostgreSQL pool connection | `gauge` | | `quickwit_metastore` | `max_connections` | Maximum number of PostgreSQL pool connections configured per metastore node | `gauge` | +Metastore gRPC clients and servers also expose shared gRPC metrics: + +| Namespace | Metric Name | Description | Labels | Type | +| --------- | ----------- | ----------- | ------ | ---- | +| `quickwit_grpc` | `requests_total` | Total number of gRPC requests processed | [`service`, `grpc_service`, `kind`, `rpc`, `status`, `code`, `metastore_kind`] | `counter` | +| `quickwit_grpc` | `requests_in_flight` | Number of gRPC requests currently in flight | [`service`, `grpc_service`, `kind`, `rpc`, `metastore_kind`] | `gauge` | +| `quickwit_grpc` | `request_duration_seconds` | Duration of gRPC requests | [`service`, `grpc_service`, `kind`, `rpc`, `status`, `code`, `metastore_kind`] | `histogram` | + +For `quickwit_grpc_requests_total`, `status="success"`, `status="error"`, and `status="cancelled"` describe the final outcome of a logical RPC. `status="retry"` describes a failed attempt that the retry policy will retry. The `code` label contains the gRPC status code for the final outcome or retry attempt. + ## Rest API Metrics | Namespace | Metric Name | Description | Type | diff --git a/quickwit/quickwit-common/src/tower/metrics.rs b/quickwit/quickwit-common/src/tower/metrics.rs index 17a9f6acec3..5449955e5f0 100644 --- a/quickwit/quickwit-common/src/tower/metrics.rs +++ b/quickwit/quickwit-common/src/tower/metrics.rs @@ -12,6 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. +use std::fmt; use std::pin::Pin; use std::task::{Context, Poll}; use std::time::Instant; @@ -19,12 +20,15 @@ use std::time::Instant; use futures::{Future, ready}; use pin_project::{pin_project, pinned_drop}; use quickwit_metrics::{ - Counter, Gauge, Histogram, Labels, LazyCounter, LazyGauge, LazyHistogram, counter, gauge, - histogram, labels, lazy_counter, lazy_gauge, lazy_histogram, + Counter, Gauge, Histogram, LabelNames, Labels, LazyCounter, LazyGauge, LazyHistogram, counter, + gauge, histogram, label_names, label_values, labels, lazy_counter, lazy_gauge, lazy_histogram, }; +use tower::retry::Policy; use tower::{Layer, Service}; +use super::retry::RetryPolicy; use crate::metrics::exponential_buckets; +use crate::retry::{RetryParams, Retryable}; pub trait RpcName { fn rpc_name() -> &'static str; @@ -69,6 +73,9 @@ fn grpc_code_label(code: tonic::Code) -> &'static str { } } +const GRPC_REQUEST_RPC_LABEL_NAME: LabelNames<1> = label_names!("rpc"); +const GRPC_REQUEST_LABEL_NAMES: LabelNames<3> = label_names!("rpc", "status", "code"); + static GRPC_REQUESTS_TOTAL: LazyCounter = lazy_counter!( name: "requests_total", description: "Total number of gRPC requests processed.", @@ -117,7 +124,7 @@ where gauge!( parent: self.requests_in_flight, - "rpc" => rpc_name, + labels: [label_values!(GRPC_REQUEST_RPC_LABEL_NAME => rpc_name)], ) .inc(); @@ -141,6 +148,36 @@ pub struct GrpcMetricsLayer { request_duration_seconds: Histogram, } +/// Retry policy that records each failed attempt which will be retried. +#[derive(Clone)] +pub struct GrpcRetryPolicy { + inner: RetryPolicy, + metrics_layer: GrpcMetricsLayer, +} + +impl Policy for GrpcRetryPolicy +where + R: Clone + RpcName, + E: fmt::Debug + GrpcStatusCode + Retryable, +{ + type Future = >::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 { + >::clone_request(&mut self.inner, request) + } +} + impl GrpcMetricsLayer { pub fn new(subsystem: &'static str, kind: &'static str) -> Self { let labels = Self::default_labels(subsystem, kind); @@ -164,6 +201,24 @@ impl GrpcMetricsLayer { } } + /// Wraps the retry policy to record retryable failed attempts as retries. + pub fn retry_policy(&self, retry_params: RetryParams) -> GrpcRetryPolicy { + GrpcRetryPolicy { + inner: RetryPolicy::from(retry_params), + metrics_layer: self.clone(), + } + } + + fn record_request(&self, rpc_name: &'static str, status: &'static str, code: tonic::Code) { + let labels = + label_values!(GRPC_REQUEST_LABEL_NAMES => rpc_name, status, grpc_code_label(code)); + counter!( + parent: self.requests_total, + labels: [labels], + ) + .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. @@ -204,11 +259,13 @@ pub struct ResponseFuture { impl PinnedDrop for ResponseFuture { fn drop(self: Pin<&mut Self>) { let elapsed = self.start.elapsed().as_secs_f64(); - let rpc_label = labels!("rpc" => self.rpc_name); - let status_label = labels!("status" => self.status); - let code_label = labels!("code" => self.code); - counter!(parent: self.requests_total, labels: [rpc_label, status_label, code_label]).inc(); - histogram!(parent: self.request_duration_seconds, labels: [rpc_label, status_label, code_label]) + let counter_labels = + label_values!(GRPC_REQUEST_LABEL_NAMES => self.rpc_name, self.status, self.code); + let histogram_labels = + label_values!(GRPC_REQUEST_LABEL_NAMES => self.rpc_name, self.status, self.code); + let rpc_label = label_values!(GRPC_REQUEST_RPC_LABEL_NAME => self.rpc_name); + counter!(parent: self.requests_total, labels: [counter_labels]).inc(); + histogram!(parent: self.request_duration_seconds, labels: [histogram_labels]) .observe(elapsed); gauge!(parent: self.requests_in_flight, labels: [rpc_label]).dec(); } @@ -240,12 +297,16 @@ where #[cfg(test)] mod tests { + use std::sync::Arc; + use std::sync::atomic::{AtomicUsize, Ordering}; + use metrics::with_local_recorder; use metrics_util::debugging::{DebugValue, DebuggingRecorder}; + use super::super::retry::RetryLayer; use super::*; - #[derive(Debug)] + #[derive(Clone, Debug)] struct HelloRequest; impl RpcName for HelloRequest { @@ -262,13 +323,47 @@ mod tests { } } + #[derive(Clone, Debug)] + struct RetrySuccessRequest; + + impl RpcName for RetrySuccessRequest { + fn rpc_name() -> &'static str { + "retry_success" + } + } + + #[derive(Clone, Debug)] + struct RetryExhaustedRequest; + + impl RpcName for RetryExhaustedRequest { + fn rpc_name() -> &'static str { + "retry_exhausted" + } + } + + #[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", @@ -301,6 +396,48 @@ 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 success_attempts = Arc::new(AtomicUsize::new(0)); + let mut retry_success_service = primary_layer.clone().layer( + RetryLayer::new(primary_layer.retry_policy(retry_params)).layer( + tower::service_fn(move |request: RetrySuccessRequest| { + let success_attempts = success_attempts.clone(); + async move { + if success_attempts.fetch_add(1, Ordering::Relaxed) < 2 { + Err(RetryableError) + } else { + Ok(request) + } + } + }), + ), + ); + retry_success_service + .call(RetrySuccessRequest) + .await + .unwrap(); + + let retry_params = RetryParams { + base_delay: std::time::Duration::ZERO, + max_delay: std::time::Duration::ZERO, + max_attempts: 3, + }; + let mut retry_exhausted_service = primary_layer.clone().layer( + RetryLayer::new(primary_layer.retry_policy(retry_params)).layer( + tower::service_fn(|_request: RetryExhaustedRequest| async move { + Err::(RetryableError) + }), + ), + ); + retry_exhausted_service + .call(RetryExhaustedRequest) + .await + .unwrap_err(); }); }); @@ -347,5 +484,21 @@ mod tests { counter_value("hello", "error", "not_found", "primary"), Some(&DebugValue::Counter(1)) ); + assert_eq!( + counter_value("retry_success", "retry", "unavailable", "primary"), + Some(&DebugValue::Counter(2)) + ); + assert_eq!( + counter_value("retry_success", "success", "ok", "primary"), + Some(&DebugValue::Counter(1)) + ); + assert_eq!( + counter_value("retry_exhausted", "retry", "unavailable", "primary"), + Some(&DebugValue::Counter(2)) + ); + assert_eq!( + counter_value("retry_exhausted", "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..21285e361fe 100644 --- a/quickwit/quickwit-serve/src/metastore.rs +++ b/quickwit/quickwit-serve/src/metastore.rs @@ -21,7 +21,7 @@ use quickwit_cluster::Cluster; use quickwit_common::pubsub::EventBroker; use quickwit_common::retry::RetryParams; use quickwit_common::tower::{ - EventListenerLayer, GrpcMetricsLayer, LoadShedLayer, RetryLayer, RetryPolicy, TimeoutLayer, + EventListenerLayer, GrpcMetricsLayer, LoadShedLayer, RetryLayer, TimeoutLayer, }; use quickwit_common::uri::Uri; use quickwit_config::service::QuickwitService; @@ -176,10 +176,13 @@ impl LocalMetastoreServer { } _ => unreachable!("unexpected metastore service `{service}`"), }; + let retry_policy = metrics_layer.retry_policy(RetryParams::standard()); Ok(MetastoreServiceClient::tower() - .stack_layer(RetryLayer::new(RetryPolicy::from(RetryParams::standard()))) - .stack_layer(TimeoutLayer::new(GRPC_METASTORE_SERVICE_TIMEOUT)) + // Keep the metrics layer outside retries: the retry policy records each retryable + // failed attempt as `retry`, while the layer records 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(), )) From 9d0eb17507c1a5d4337ed64484c408468e3d50e7 Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Mon, 17 Aug 2026 18:46:58 +0200 Subject: [PATCH 08/12] Remove gRPC metrics documentation changes Keep this change focused on retry metric behavior without expanding the metrics reference. --- docs/reference/metrics.md | 10 ---------- 1 file changed, 10 deletions(-) diff --git a/docs/reference/metrics.md b/docs/reference/metrics.md index 293903fcb76..33a49854895 100644 --- a/docs/reference/metrics.md +++ b/docs/reference/metrics.md @@ -69,16 +69,6 @@ PostgreSQL-backed metastores also expose connection pool gauges: | `quickwit_metastore` | `acquire_connections` | Number of requests currently waiting to acquire a PostgreSQL pool connection | `gauge` | | `quickwit_metastore` | `max_connections` | Maximum number of PostgreSQL pool connections configured per metastore node | `gauge` | -Metastore gRPC clients and servers also expose shared gRPC metrics: - -| Namespace | Metric Name | Description | Labels | Type | -| --------- | ----------- | ----------- | ------ | ---- | -| `quickwit_grpc` | `requests_total` | Total number of gRPC requests processed | [`service`, `grpc_service`, `kind`, `rpc`, `status`, `code`, `metastore_kind`] | `counter` | -| `quickwit_grpc` | `requests_in_flight` | Number of gRPC requests currently in flight | [`service`, `grpc_service`, `kind`, `rpc`, `metastore_kind`] | `gauge` | -| `quickwit_grpc` | `request_duration_seconds` | Duration of gRPC requests | [`service`, `grpc_service`, `kind`, `rpc`, `status`, `code`, `metastore_kind`] | `histogram` | - -For `quickwit_grpc_requests_total`, `status="success"`, `status="error"`, and `status="cancelled"` describe the final outcome of a logical RPC. `status="retry"` describes a failed attempt that the retry policy will retry. The `code` label contains the gRPC status code for the final outcome or retry attempt. - ## Rest API Metrics | Namespace | Metric Name | Description | Type | From 0d434be5dee1a059213ead623e638e321b08461f Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Mon, 17 Aug 2026 18:55:39 +0200 Subject: [PATCH 09/12] Clarify metastore retry metric layering Document why request metrics must wrap retries to avoid counting retry attempts as final errors. --- quickwit/quickwit-serve/src/metastore.rs | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/quickwit/quickwit-serve/src/metastore.rs b/quickwit/quickwit-serve/src/metastore.rs index 21285e361fe..d52198767b2 100644 --- a/quickwit/quickwit-serve/src/metastore.rs +++ b/quickwit/quickwit-serve/src/metastore.rs @@ -178,8 +178,9 @@ impl LocalMetastoreServer { }; let retry_policy = metrics_layer.retry_policy(RetryParams::standard()); Ok(MetastoreServiceClient::tower() - // Keep the metrics layer outside retries: the retry policy records each retryable - // failed attempt as `retry`, while the layer records the final outcome. + // Keep metrics outside retries so the response future records only the final outcome + // seen by the caller. The retry policy records attempts that will be retried as + // `status="retry"`; placing metrics inside retries would also count them as errors. .stack_layer(metrics_layer) .stack_layer(RetryLayer::new(retry_policy)) .stack_layer(TimeoutLayer::new(GRPC_METASTORE_SERVICE_TIMEOUT)) From bb256e0cbb01403e9aa934a801f7cfaf516aac6a Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Mon, 17 Aug 2026 19:00:31 +0200 Subject: [PATCH 10/12] Explain metastore retry layer order Document Tower's outermost-first stacking and why metrics must wrap retries. --- quickwit/quickwit-serve/src/metastore.rs | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/quickwit/quickwit-serve/src/metastore.rs b/quickwit/quickwit-serve/src/metastore.rs index d52198767b2..e5bb6cb583e 100644 --- a/quickwit/quickwit-serve/src/metastore.rs +++ b/quickwit/quickwit-serve/src/metastore.rs @@ -178,9 +178,10 @@ impl LocalMetastoreServer { }; let retry_policy = metrics_layer.retry_policy(RetryParams::standard()); Ok(MetastoreServiceClient::tower() - // Keep metrics outside retries so the response future records only the final outcome - // seen by the caller. The retry policy records attempts that will be retried as - // `status="retry"`; placing metrics inside retries would also count them as errors. + // Layer order matters: the first `stack_layer` is the outermost layer. Metrics must + // wrap retries so the response future records only the final outcome seen by the + // caller, while the retry policy records attempts as `status="retry"`. Reversing + // these layers would also record every retryable failed attempt as `status="error"`. .stack_layer(metrics_layer) .stack_layer(RetryLayer::new(retry_policy)) .stack_layer(TimeoutLayer::new(GRPC_METASTORE_SERVICE_TIMEOUT)) From 6b74858fd8a46bcc53ab33379de4ad643f4e631c Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Tue, 18 Aug 2026 09:59:55 +0200 Subject: [PATCH 11/12] Integrate generic retry policy changes Decouple gRPC retry metrics from Quickwit's concrete retry policy while preserving retry outcome coverage. --- quickwit/quickwit-common/src/tower/metrics.rs | 71 ++++++++++--------- quickwit/quickwit-serve/src/metastore.rs | 10 ++- 2 files changed, 41 insertions(+), 40 deletions(-) diff --git a/quickwit/quickwit-common/src/tower/metrics.rs b/quickwit/quickwit-common/src/tower/metrics.rs index 5449955e5f0..03865621e94 100644 --- a/quickwit/quickwit-common/src/tower/metrics.rs +++ b/quickwit/quickwit-common/src/tower/metrics.rs @@ -12,7 +12,6 @@ // See the License for the specific language governing permissions and // limitations under the License. -use std::fmt; use std::pin::Pin; use std::task::{Context, Poll}; use std::time::Instant; @@ -26,9 +25,7 @@ use quickwit_metrics::{ use tower::retry::Policy; use tower::{Layer, Service}; -use super::retry::RetryPolicy; use crate::metrics::exponential_buckets; -use crate::retry::{RetryParams, Retryable}; pub trait RpcName { fn rpc_name() -> &'static str; @@ -150,17 +147,18 @@ pub struct GrpcMetricsLayer { /// Retry policy that records each failed attempt which will be retried. #[derive(Clone)] -pub struct GrpcRetryPolicy { - inner: RetryPolicy, +pub struct GrpcRetryPolicy

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

where - R: Clone + RpcName, - E: fmt::Debug + GrpcStatusCode + Retryable, + P: Policy, + R: RpcName, + E: GrpcStatusCode, { - type Future = >::Future; + type Future = P::Future; fn retry(&mut self, request: &mut R, result: &mut Result) -> Option { let retry = self.inner.retry(request, result); @@ -174,7 +172,7 @@ where } fn clone_request(&mut self, request: &R) -> Option { - >::clone_request(&mut self.inner, request) + self.inner.clone_request(request) } } @@ -201,10 +199,10 @@ impl GrpcMetricsLayer { } } - /// Wraps the retry policy to record retryable failed attempts as retries. - pub fn retry_policy(&self, retry_params: RetryParams) -> GrpcRetryPolicy { + /// Wraps a retry policy to record failed attempts that it retries. + pub fn from_retry_policy

(&self, inner: P) -> GrpcRetryPolicy

{ GrpcRetryPolicy { - inner: RetryPolicy::from(retry_params), + inner, metrics_layer: self.clone(), } } @@ -303,8 +301,9 @@ mod tests { use metrics::with_local_recorder; use metrics_util::debugging::{DebugValue, DebuggingRecorder}; - use super::super::retry::RetryLayer; + use super::super::retry::{RetryLayer, RetryPolicy}; use super::*; + use crate::retry::{RetryParams, Retryable}; #[derive(Clone, Debug)] struct HelloRequest; @@ -403,20 +402,22 @@ mod tests { max_attempts: 3, }; let success_attempts = Arc::new(AtomicUsize::new(0)); - let mut retry_success_service = primary_layer.clone().layer( - RetryLayer::new(primary_layer.retry_policy(retry_params)).layer( - tower::service_fn(move |request: RetrySuccessRequest| { - let success_attempts = success_attempts.clone(); - async move { - if success_attempts.fetch_add(1, Ordering::Relaxed) < 2 { - Err(RetryableError) - } else { - Ok(request) + let retry_policy = primary_layer.from_retry_policy(RetryPolicy::from(retry_params)); + let mut retry_success_service = + primary_layer + .clone() + .layer(RetryLayer::new(retry_policy).layer(tower::service_fn( + move |request: RetrySuccessRequest| { + let success_attempts = success_attempts.clone(); + async move { + if success_attempts.fetch_add(1, Ordering::Relaxed) < 2 { + Err(RetryableError) + } else { + Ok(request) + } } - } - }), - ), - ); + }, + ))); retry_success_service .call(RetrySuccessRequest) .await @@ -427,13 +428,15 @@ mod tests { max_delay: std::time::Duration::ZERO, max_attempts: 3, }; - let mut retry_exhausted_service = primary_layer.clone().layer( - RetryLayer::new(primary_layer.retry_policy(retry_params)).layer( - tower::service_fn(|_request: RetryExhaustedRequest| async move { - Err::(RetryableError) - }), - ), - ); + let retry_policy = primary_layer.from_retry_policy(RetryPolicy::from(retry_params)); + let mut retry_exhausted_service = + primary_layer + .clone() + .layer(RetryLayer::new(retry_policy).layer(tower::service_fn( + |_request: RetryExhaustedRequest| async move { + Err::(RetryableError) + }, + ))); retry_exhausted_service .call(RetryExhaustedRequest) .await diff --git a/quickwit/quickwit-serve/src/metastore.rs b/quickwit/quickwit-serve/src/metastore.rs index e5bb6cb583e..f42dc80dbf5 100644 --- a/quickwit/quickwit-serve/src/metastore.rs +++ b/quickwit/quickwit-serve/src/metastore.rs @@ -21,7 +21,7 @@ use quickwit_cluster::Cluster; use quickwit_common::pubsub::EventBroker; use quickwit_common::retry::RetryParams; use quickwit_common::tower::{ - EventListenerLayer, GrpcMetricsLayer, LoadShedLayer, RetryLayer, TimeoutLayer, + EventListenerLayer, GrpcMetricsLayer, LoadShedLayer, RetryLayer, RetryPolicy, TimeoutLayer, }; use quickwit_common::uri::Uri; use quickwit_config::service::QuickwitService; @@ -176,12 +176,10 @@ impl LocalMetastoreServer { } _ => unreachable!("unexpected metastore service `{service}`"), }; - let retry_policy = metrics_layer.retry_policy(RetryParams::standard()); + let retry_policy = + metrics_layer.from_retry_policy(RetryPolicy::from(RetryParams::standard())); Ok(MetastoreServiceClient::tower() - // Layer order matters: the first `stack_layer` is the outermost layer. Metrics must - // wrap retries so the response future records only the final outcome seen by the - // caller, while the retry policy records attempts as `status="retry"`. Reversing - // these layers would also record every retryable failed attempt as `status="error"`. + // 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)) From dce61d205b050a8d057e71e072cc7305f9380d67 Mon Sep 17 00:00:00 2001 From: Luca Cominardi Date: Tue, 18 Aug 2026 10:12:01 +0200 Subject: [PATCH 12/12] Adopt alternative retry metrics implementation Align the metrics wrapper and tests with the proposed generic policy implementation. --- quickwit/quickwit-common/src/tower/metrics.rs | 114 ++++-------------- 1 file changed, 25 insertions(+), 89 deletions(-) diff --git a/quickwit/quickwit-common/src/tower/metrics.rs b/quickwit/quickwit-common/src/tower/metrics.rs index 03865621e94..cb3cc9884e3 100644 --- a/quickwit/quickwit-common/src/tower/metrics.rs +++ b/quickwit/quickwit-common/src/tower/metrics.rs @@ -19,8 +19,8 @@ use std::time::Instant; use futures::{Future, ready}; use pin_project::{pin_project, pinned_drop}; use quickwit_metrics::{ - Counter, Gauge, Histogram, LabelNames, Labels, LazyCounter, LazyGauge, LazyHistogram, counter, - gauge, histogram, label_names, label_values, labels, lazy_counter, lazy_gauge, lazy_histogram, + 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}; @@ -70,9 +70,6 @@ fn grpc_code_label(code: tonic::Code) -> &'static str { } } -const GRPC_REQUEST_RPC_LABEL_NAME: LabelNames<1> = label_names!("rpc"); -const GRPC_REQUEST_LABEL_NAMES: LabelNames<3> = label_names!("rpc", "status", "code"); - static GRPC_REQUESTS_TOTAL: LazyCounter = lazy_counter!( name: "requests_total", description: "Total number of gRPC requests processed.", @@ -121,7 +118,7 @@ where gauge!( parent: self.requests_in_flight, - labels: [label_values!(GRPC_REQUEST_RPC_LABEL_NAME => rpc_name)], + "rpc" => rpc_name, ) .inc(); @@ -208,11 +205,13 @@ impl GrpcMetricsLayer { } fn record_request(&self, rpc_name: &'static str, status: &'static str, code: tonic::Code) { - let labels = - label_values!(GRPC_REQUEST_LABEL_NAMES => rpc_name, status, grpc_code_label(code)); counter!( parent: self.requests_total, - labels: [labels], + labels: [labels!( + "rpc" => rpc_name, + "status" => status, + "code" => grpc_code_label(code), + )], ) .inc(); } @@ -257,13 +256,11 @@ pub struct ResponseFuture { impl PinnedDrop for ResponseFuture { fn drop(self: Pin<&mut Self>) { let elapsed = self.start.elapsed().as_secs_f64(); - let counter_labels = - label_values!(GRPC_REQUEST_LABEL_NAMES => self.rpc_name, self.status, self.code); - let histogram_labels = - label_values!(GRPC_REQUEST_LABEL_NAMES => self.rpc_name, self.status, self.code); - let rpc_label = label_values!(GRPC_REQUEST_RPC_LABEL_NAME => self.rpc_name); - counter!(parent: self.requests_total, labels: [counter_labels]).inc(); - histogram!(parent: self.request_duration_seconds, labels: [histogram_labels]) + let rpc_label = labels!("rpc" => self.rpc_name); + let status_label = labels!("status" => self.status); + let code_label = labels!("code" => self.code); + counter!(parent: self.requests_total, labels: [rpc_label, status_label, code_label]).inc(); + histogram!(parent: self.request_duration_seconds, labels: [rpc_label, status_label, code_label]) .observe(elapsed); gauge!(parent: self.requests_in_flight, labels: [rpc_label]).dec(); } @@ -295,9 +292,6 @@ where #[cfg(test)] mod tests { - use std::sync::Arc; - use std::sync::atomic::{AtomicUsize, Ordering}; - use metrics::with_local_recorder; use metrics_util::debugging::{DebugValue, DebuggingRecorder}; @@ -322,24 +316,6 @@ mod tests { } } - #[derive(Clone, Debug)] - struct RetrySuccessRequest; - - impl RpcName for RetrySuccessRequest { - fn rpc_name() -> &'static str { - "retry_success" - } - } - - #[derive(Clone, Debug)] - struct RetryExhaustedRequest; - - impl RpcName for RetryExhaustedRequest { - fn rpc_name() -> &'static str { - "retry_exhausted" - } - } - #[derive(Debug)] struct RetryableError; @@ -383,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(); @@ -401,46 +378,13 @@ mod tests { max_delay: std::time::Duration::ZERO, max_attempts: 3, }; - let success_attempts = Arc::new(AtomicUsize::new(0)); let retry_policy = primary_layer.from_retry_policy(RetryPolicy::from(retry_params)); - let mut retry_success_service = - primary_layer - .clone() - .layer(RetryLayer::new(retry_policy).layer(tower::service_fn( - move |request: RetrySuccessRequest| { - let success_attempts = success_attempts.clone(); - async move { - if success_attempts.fetch_add(1, Ordering::Relaxed) < 2 { - Err(RetryableError) - } else { - Ok(request) - } - } - }, - ))); - retry_success_service - .call(RetrySuccessRequest) - .await - .unwrap(); - - 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_exhausted_service = - primary_layer - .clone() - .layer(RetryLayer::new(retry_policy).layer(tower::service_fn( - |_request: RetryExhaustedRequest| async move { - Err::(RetryableError) - }, - ))); - retry_exhausted_service - .call(RetryExhaustedRequest) - .await - .unwrap_err(); + 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(); }); }); @@ -488,19 +432,11 @@ mod tests { Some(&DebugValue::Counter(1)) ); assert_eq!( - counter_value("retry_success", "retry", "unavailable", "primary"), - Some(&DebugValue::Counter(2)) - ); - assert_eq!( - counter_value("retry_success", "success", "ok", "primary"), - Some(&DebugValue::Counter(1)) - ); - assert_eq!( - counter_value("retry_exhausted", "retry", "unavailable", "primary"), + counter_value("hello", "retry", "unavailable", "primary"), Some(&DebugValue::Counter(2)) ); assert_eq!( - counter_value("retry_exhausted", "error", "unavailable", "primary"), + counter_value("hello", "error", "unavailable", "primary"), Some(&DebugValue::Counter(1)) ); }