diff --git a/quickwit/quickwit-common/src/tower/metrics.rs b/quickwit/quickwit-common/src/tower/metrics.rs index 17a9f6acec3..cb3cc9884e3 100644 --- a/quickwit/quickwit-common/src/tower/metrics.rs +++ b/quickwit/quickwit-common/src/tower/metrics.rs @@ -22,6 +22,7 @@ use quickwit_metrics::{ Counter, Gauge, Histogram, Labels, LazyCounter, LazyGauge, LazyHistogram, counter, gauge, histogram, labels, lazy_counter, lazy_gauge, lazy_histogram, }; +use tower::retry::Policy; use tower::{Layer, Service}; use crate::metrics::exponential_buckets; @@ -141,6 +142,37 @@ pub struct GrpcMetricsLayer { request_duration_seconds: Histogram, } +/// Retry policy that records each failed attempt which will be retried. +#[derive(Clone)] +pub struct GrpcRetryPolicy
{
+ inner: P,
+ metrics_layer: GrpcMetricsLayer,
+}
+
+impl
+where
+ P: Policy (&self, inner: P) -> GrpcRetryPolicy {
+ GrpcRetryPolicy {
+ inner,
+ metrics_layer: self.clone(),
+ }
+ }
+
+ fn record_request(&self, rpc_name: &'static str, status: &'static str, code: tonic::Code) {
+ counter!(
+ parent: self.requests_total,
+ labels: [labels!(
+ "rpc" => rpc_name,
+ "status" => status,
+ "code" => grpc_code_label(code),
+ )],
+ )
+ .inc();
+ }
+
fn default_labels(subsystem: &'static str, kind: &'static str) -> Labels<3> {
// `service` is kept for backward compatibility with existing consumers. Prefer
// `grpc_service` for new consumers.
@@ -243,9 +295,11 @@ mod tests {
use metrics::with_local_recorder;
use metrics_util::debugging::{DebugValue, DebuggingRecorder};
+ use super::super::retry::{RetryLayer, RetryPolicy};
use super::*;
+ use crate::retry::{RetryParams, Retryable};
- #[derive(Debug)]
+ #[derive(Clone, Debug)]
struct HelloRequest;
impl RpcName for HelloRequest {
@@ -262,13 +316,29 @@ mod tests {
}
}
+ #[derive(Debug)]
+ struct RetryableError;
+
+ impl GrpcStatusCode for RetryableError {
+ fn grpc_status_code(&self) -> tonic::Code {
+ tonic::Code::Unavailable
+ }
+ }
+
+ impl Retryable for RetryableError {
+ fn is_retryable(&self) -> bool {
+ true
+ }
+ }
+
#[test]
fn test_grpc_metrics() {
let recorder = DebuggingRecorder::new();
let snapshotter = recorder.snapshotter();
with_local_recorder(&recorder, || {
- futures::executor::block_on(async {
+ let runtime = tokio::runtime::Runtime::new().unwrap();
+ runtime.block_on(async {
let primary_layer = GrpcMetricsLayer::new_with_labels(
"quickwit_test",
"server",
@@ -289,10 +359,11 @@ mod tests {
let mut read_replica_service = read_replica_layer.layer(tower::service_fn(
|request: HelloRequest| async move { Ok::<_, tonic::Status>(request) },
));
- let mut failing_service =
- primary_layer.layer(tower::service_fn(|_request: HelloRequest| async move {
+ let mut failing_service = primary_layer.clone().layer(tower::service_fn(
+ |_request: HelloRequest| async move {
Err::