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::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 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 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#[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 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}