1use 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
54fn 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 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 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#[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 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}