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::collector::metrics::v1::ExportMetricsServiceRequest;
28use otel_arrow_rust::proto::opentelemetry::metrics::v1::metric;
29use session::protocol_ctx::{OtlpMetricCtx, ProtocolCtx};
30use tonic::metadata::{Entry, MetadataValue};
31use tonic::service::Interceptor;
32use tonic::{Request, Response, Status, Streaming};
33
34use crate::error;
35use crate::grpc::context_auth;
36use crate::query_handler::{MetricsIngestOutcome, OpenTelemetryProtocolHandlerRef};
37
38const EXPONENTIAL_HISTOGRAM_UNSUPPORTED: &str = "OTel Arrow exponential histograms are unsupported because the Arrow wire format omits zero_threshold";
39
40pub struct OtelArrowServiceHandler<T> {
41    handler: T,
42    user_provider: Option<UserProviderRef>,
43}
44
45impl<T> OtelArrowServiceHandler<T> {
46    pub fn new(handler: T, user_provider: Option<UserProviderRef>) -> Self {
47        Self {
48            handler,
49            user_provider,
50        }
51    }
52}
53
54/// Removes unsupported Arrow histograms before ingestion, preserving other metrics.
55fn remove_exponential_histograms(request: &mut ExportMetricsServiceRequest) -> bool {
56    let mut has_data_points = false;
57    for scope in request
58        .resource_metrics
59        .iter_mut()
60        .flat_map(|resource| &mut resource.scope_metrics)
61    {
62        scope.metrics.retain(|item| {
63            if let Some(metric::Data::ExponentialHistogram(histogram)) = &item.data {
64                has_data_points |= !histogram.data_points.is_empty();
65                false
66            } else {
67                true
68            }
69        });
70    }
71    has_data_points
72}
73
74fn batch_status(
75    batch_id: i64,
76    outcome: MetricsIngestOutcome,
77    has_exponential_histogram_data_points: bool,
78) -> BatchStatus {
79    let status_code = if outcome.accepted_data_points == 0
80        && (outcome.rejected_data_points > 0 || has_exponential_histogram_data_points)
81    {
82        ArrowStatusCode::InvalidArgument
83    } else {
84        ArrowStatusCode::Ok
85    };
86    let status_message = if has_exponential_histogram_data_points {
87        EXPONENTIAL_HISTOGRAM_UNSUPPORTED.to_string()
88    } else {
89        outcome.error_message.unwrap_or_default()
90    };
91    BatchStatus {
92        batch_id,
93        status_code: status_code as i32,
94        status_message,
95    }
96}
97
98#[async_trait::async_trait]
99impl ArrowMetricsService for OtelArrowServiceHandler<OpenTelemetryProtocolHandlerRef> {
100    type ArrowMetricsStream = futures::channel::mpsc::Receiver<Result<BatchStatus, Status>>;
101    async fn arrow_metrics(
102        &self,
103        request: Request<Streaming<BatchArrowRecords>>,
104    ) -> Result<Response<Self::ArrowMetricsStream>, Status> {
105        let (mut sender, receiver) = futures::channel::mpsc::channel(100);
106
107        let (headers, extensions, mut incoming_requests) = request.into_parts();
108
109        // OTEL Arrow is currently used only for external ingestion. Its service bypasses
110        // the frontend router middleware, so even the internal listener defaults to
111        // Channel::Grpc and remains rate-limited. Before using this path internally,
112        // propagate the server-owned Channel::Internal marker to this service.
113        let query_ctx =
114            context_auth::create_query_context_from_grpc_metadata(&headers, &extensions)?;
115        context_auth::check_auth(self.user_provider.clone(), &headers, query_ctx.clone()).await?;
116        let query_ctx = {
117            let mut ctx = query_ctx.fork();
118            ctx.set_protocol_ctx(ProtocolCtx::OtlpMetric(OtlpMetricCtx::default()));
119            Arc::new(ctx)
120        };
121
122        let handler = self.handler.clone();
123
124        // handles incoming requests
125        common_runtime::spawn_global(async move {
126            let mut consumer = Consumer::default();
127            while let Some(batch_res) = incoming_requests.message().await.transpose() {
128                let mut batch = match batch_res {
129                    Ok(batch) => batch,
130                    Err(e) => {
131                        error!(
132                            "Failed to receive batch from otel-arrow client, error: {}",
133                            e
134                        );
135                        let _ = sender.send(Err(e)).await;
136                        return;
137                    }
138                };
139                let batch_id = batch.batch_id;
140                let mut request = match consumer.consume_metrics_batches(&mut batch).map_err(|e| {
141                    error::HandleOtelArrowRequestSnafu {
142                        err_msg: e.to_string(),
143                    }
144                    .build()
145                }) {
146                    Ok(request) => request,
147                    Err(e) => {
148                        let _ = sender
149                            .send(Err(Status::new(
150                                status_to_tonic_code(e.status_code()),
151                                e.to_string(),
152                            )))
153                            .await;
154                        error!(e;
155                            "Failed to consume batch from otel-arrow client"
156                        );
157                        return;
158                    }
159                };
160                let has_exponential_histogram_data_points =
161                    remove_exponential_histograms(&mut request);
162                let outcome = match handler.metrics(request, query_ctx.clone()).await {
163                    Ok(outcome) => outcome,
164                    Err(error::Error::InvalidOtlpMetricInput { reason }) => {
165                        let _ = sender
166                            .send(Ok(BatchStatus {
167                                batch_id,
168                                status_code: ArrowStatusCode::InvalidArgument as i32,
169                                status_message: reason,
170                            }))
171                            .await;
172                        continue;
173                    }
174                    Err(e) => {
175                        let _ = sender
176                            .send(Err(Status::new(
177                                status_to_tonic_code(e.status_code()),
178                                e.to_string(),
179                            )))
180                            .await;
181                        error!(e; "Failed to ingest metrics from otel-arrow");
182                        return;
183                    }
184                };
185                let batch_status =
186                    batch_status(batch_id, outcome, has_exponential_histogram_data_points);
187                let _ = sender.send(Ok(batch_status)).await;
188            }
189        });
190        Ok(Response::new(receiver))
191    }
192}
193
194/// This serves as a workaround for otel-arrow collector's custom header.
195#[derive(Clone)]
196pub struct HeaderInterceptor;
197
198impl Interceptor for HeaderInterceptor {
199    fn call(&mut self, mut request: Request<()>) -> Result<Request<()>, Status> {
200        if let Ok(Entry::Occupied(mut e)) = request.metadata_mut().entry("grpc-encoding") {
201            // This works as a workaround to handle customized compression type (zstdarrow*) in otel-arrow.
202            if e.get().as_bytes().starts_with(b"zstdarrow") {
203                e.insert(MetadataValue::from_static("zstd"));
204            }
205        }
206        Ok(request)
207    }
208}
209
210#[cfg(test)]
211mod tests {
212    use super::*;
213
214    #[test]
215    fn removes_arrow_exponential_histograms_and_preserves_other_metrics() {
216        use otel_arrow_rust::proto::opentelemetry::metrics::v1::{
217            ExponentialHistogram, ExponentialHistogramDataPoint, Gauge, Metric, ResourceMetrics,
218            ScopeMetrics,
219        };
220
221        let mut request = ExportMetricsServiceRequest {
222            resource_metrics: vec![ResourceMetrics {
223                scope_metrics: vec![ScopeMetrics {
224                    metrics: vec![
225                        Metric {
226                            data: Some(metric::Data::ExponentialHistogram(ExponentialHistogram {
227                                data_points: vec![ExponentialHistogramDataPoint::default()],
228                                ..Default::default()
229                            })),
230                            ..Default::default()
231                        },
232                        Metric {
233                            data: Some(metric::Data::Gauge(Gauge::default())),
234                            ..Default::default()
235                        },
236                    ],
237                    ..Default::default()
238                }],
239                ..Default::default()
240            }],
241        };
242        assert!(remove_exponential_histograms(&mut request));
243        let metrics = &request.resource_metrics[0].scope_metrics[0].metrics;
244        assert_eq!(metrics.len(), 1);
245        assert!(matches!(metrics[0].data, Some(metric::Data::Gauge(_))));
246        assert!(!remove_exponential_histograms(&mut request));
247    }
248
249    #[test]
250    fn batch_status_explains_arrow_exponential_histogram_limit() {
251        for accepted_data_points in [0, 1] {
252            let status = batch_status(
253                7,
254                MetricsIngestOutcome {
255                    accepted_data_points,
256                    ..Default::default()
257                },
258                true,
259            );
260            assert_eq!(7, status.batch_id);
261            assert_eq!(
262                if accepted_data_points == 0 {
263                    ArrowStatusCode::InvalidArgument
264                } else {
265                    ArrowStatusCode::Ok
266                } as i32,
267                status.status_code
268            );
269            assert_eq!(EXPONENTIAL_HISTOGRAM_UNSUPPORTED, status.status_message);
270        }
271    }
272}