Skip to main content

servers/
otel_arrow.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::sync::Arc;
16
17use auth::UserProviderRef;
18use common_error::ext::ErrorExt;
19use common_error::status_code::status_to_tonic_code;
20use common_telemetry::error;
21use futures::SinkExt;
22use otel_arrow_rust::Consumer;
23use otel_arrow_rust::proto::opentelemetry::arrow::v1::arrow_metrics_service_server::ArrowMetricsService;
24use otel_arrow_rust::proto::opentelemetry::arrow::v1::{
25    BatchArrowRecords, BatchStatus, StatusCode as ArrowStatusCode,
26};
27use otel_arrow_rust::proto::opentelemetry::metrics::v1::metric;
28use session::protocol_ctx::{OtlpMetricCtx, ProtocolCtx};
29use tonic::metadata::{Entry, MetadataValue};
30use tonic::service::Interceptor;
31use tonic::{Request, Response, Status, Streaming};
32
33use crate::error;
34use crate::grpc::context_auth;
35use crate::query_handler::{MetricsIngestOutcome, OpenTelemetryProtocolHandlerRef};
36
37const EXPONENTIAL_HISTOGRAM_UNSUPPORTED: &str = "OTel Arrow exponential histograms are unsupported because the Arrow wire format omits zero_threshold";
38
39pub struct OtelArrowServiceHandler<T> {
40    handler: T,
41    user_provider: Option<UserProviderRef>,
42}
43
44impl<T> OtelArrowServiceHandler<T> {
45    pub fn new(handler: T, user_provider: Option<UserProviderRef>) -> Self {
46        Self {
47            handler,
48            user_provider,
49        }
50    }
51}
52
53fn batch_status(
54    batch_id: i64,
55    outcome: MetricsIngestOutcome,
56    has_exponential_histogram_data_points: bool,
57) -> BatchStatus {
58    let status_code = if outcome.accepted_data_points == 0 && outcome.rejected_data_points > 0 {
59        ArrowStatusCode::InvalidArgument
60    } else {
61        ArrowStatusCode::Ok
62    };
63    let status_message = match outcome.error_message {
64        // Arrow keeps the feature gate off, so these fail before per-point validation.
65        Some(_) if has_exponential_histogram_data_points => {
66            EXPONENTIAL_HISTOGRAM_UNSUPPORTED.to_string()
67        }
68        Some(message) => message,
69        None => String::new(),
70    };
71    BatchStatus {
72        batch_id,
73        status_code: status_code as i32,
74        status_message,
75    }
76}
77
78#[async_trait::async_trait]
79impl ArrowMetricsService for OtelArrowServiceHandler<OpenTelemetryProtocolHandlerRef> {
80    type ArrowMetricsStream = futures::channel::mpsc::Receiver<Result<BatchStatus, Status>>;
81    async fn arrow_metrics(
82        &self,
83        request: Request<Streaming<BatchArrowRecords>>,
84    ) -> Result<Response<Self::ArrowMetricsStream>, Status> {
85        let (mut sender, receiver) = futures::channel::mpsc::channel(100);
86
87        let (headers, _, mut incoming_requests) = request.into_parts();
88
89        let query_ctx = context_auth::create_query_context_from_grpc_metadata(&headers)?;
90        context_auth::check_auth(self.user_provider.clone(), &headers, query_ctx.clone()).await?;
91        let query_ctx = {
92            let mut ctx = query_ctx.fork();
93            ctx.set_protocol_ctx(ProtocolCtx::OtlpMetric(OtlpMetricCtx::default()));
94            Arc::new(ctx)
95        };
96
97        let handler = self.handler.clone();
98
99        // handles incoming requests
100        common_runtime::spawn_global(async move {
101            let mut consumer = Consumer::default();
102            while let Some(batch_res) = incoming_requests.message().await.transpose() {
103                let mut batch = match batch_res {
104                    Ok(batch) => batch,
105                    Err(e) => {
106                        error!(
107                            "Failed to receive batch from otel-arrow client, error: {}",
108                            e
109                        );
110                        let _ = sender.send(Err(e)).await;
111                        return;
112                    }
113                };
114                let batch_id = batch.batch_id;
115                let request = match consumer.consume_metrics_batches(&mut batch).map_err(|e| {
116                    error::HandleOtelArrowRequestSnafu {
117                        err_msg: e.to_string(),
118                    }
119                    .build()
120                }) {
121                    Ok(request) => request,
122                    Err(e) => {
123                        let _ = sender
124                            .send(Err(Status::new(
125                                status_to_tonic_code(e.status_code()),
126                                e.to_string(),
127                            )))
128                            .await;
129                        error!(e;
130                            "Failed to consume batch from otel-arrow client"
131                        );
132                        return;
133                    }
134                };
135                let has_exponential_histogram_data_points = request
136                    .resource_metrics
137                    .iter()
138                    .flat_map(|resource| &resource.scope_metrics)
139                    .flat_map(|scope| &scope.metrics)
140                    .any(|item| {
141                        matches!(
142                            item.data.as_ref(),
143                            Some(metric::Data::ExponentialHistogram(histogram))
144                                if !histogram.data_points.is_empty()
145                        )
146                    });
147                let outcome = match handler.metrics(request, query_ctx.clone()).await {
148                    Ok(outcome) => outcome,
149                    Err(error::Error::InvalidOtlpMetricInput { reason }) => {
150                        let _ = sender
151                            .send(Ok(BatchStatus {
152                                batch_id,
153                                status_code: ArrowStatusCode::InvalidArgument as i32,
154                                status_message: reason,
155                            }))
156                            .await;
157                        continue;
158                    }
159                    Err(e) => {
160                        let _ = sender
161                            .send(Err(Status::new(
162                                status_to_tonic_code(e.status_code()),
163                                e.to_string(),
164                            )))
165                            .await;
166                        error!(e; "Failed to ingest metrics from otel-arrow");
167                        return;
168                    }
169                };
170                let batch_status =
171                    batch_status(batch_id, outcome, has_exponential_histogram_data_points);
172                let _ = sender.send(Ok(batch_status)).await;
173            }
174        });
175        Ok(Response::new(receiver))
176    }
177}
178
179/// This serves as a workaround for otel-arrow collector's custom header.
180#[derive(Clone)]
181pub struct HeaderInterceptor;
182
183impl Interceptor for HeaderInterceptor {
184    fn call(&mut self, mut request: Request<()>) -> Result<Request<()>, Status> {
185        if let Ok(Entry::Occupied(mut e)) = request.metadata_mut().entry("grpc-encoding") {
186            // This works as a workaround to handle customized compression type (zstdarrow*) in otel-arrow.
187            if e.get().as_bytes().starts_with(b"zstdarrow") {
188                e.insert(MetadataValue::from_static("zstd"));
189            }
190        }
191        Ok(request)
192    }
193}
194
195#[cfg(test)]
196mod tests {
197    use super::*;
198
199    #[test]
200    fn batch_status_explains_arrow_exponential_histogram_limit() {
201        let status = batch_status(
202            7,
203            MetricsIngestOutcome {
204                rejected_data_points: 1,
205                error_message: Some("internal OTLP rejection detail".to_string()),
206                ..Default::default()
207            },
208            true,
209        );
210
211        assert_eq!(7, status.batch_id);
212        assert_eq!(ArrowStatusCode::InvalidArgument as i32, status.status_code);
213        assert_eq!(EXPONENTIAL_HISTOGRAM_UNSUPPORTED, status.status_message);
214    }
215}