1use std::sync::Arc;
16
17use axum::Extension;
18use axum::extract::State;
19use axum::http::{StatusCode, header};
20use axum::response::IntoResponse;
21use axum_extra::TypedHeader;
22use bytes::Bytes;
23use common_catalog::consts::{TRACE_TABLE_NAME, TRACE_TABLE_NAME_SESSION_KEY};
24use common_telemetry::tracing;
25use headers::ContentType;
26use mime_guess::mime;
27use opentelemetry_proto::tonic::collector::logs::v1::{
28 ExportLogsServiceRequest, ExportLogsServiceResponse,
29};
30use opentelemetry_proto::tonic::collector::metrics::v1::{
31 ExportMetricsPartialSuccess, ExportMetricsServiceResponse,
32};
33use opentelemetry_proto::tonic::collector::trace::v1::{
34 ExportTracePartialSuccess, ExportTraceServiceRequest, ExportTraceServiceResponse,
35};
36use otel_arrow_rust::proto::opentelemetry::collector::metrics::v1::ExportMetricsServiceRequest;
37use pipeline::PipelineWay;
38use prost::Message;
39use session::context::{Channel, QueryContext};
40use session::protocol_ctx::{MetricType, OtlpMetricCtx, ProtocolCtx};
41use snafu::prelude::*;
42
43use crate::error::{self, PipelineSnafu, Result};
44use crate::http::extractor::{
45 LogTableName, OtlpMetricOptions, PipelineInfo, SelectInfoWrapper, TraceTableName,
46};
47use crate::http::header::{CONTENT_TYPE_PROTOBUF, write_cost_header_map};
48use crate::metrics::METRIC_HTTP_OPENTELEMETRY_LOGS_ELAPSED;
49use crate::query_handler::{
50 MetricsIngestOutcome, OpenTelemetryProtocolHandlerRef, PipelineHandler, TraceIngestOutcome,
51};
52
53#[derive(Clone, prost::Message)]
54pub struct GoogleRpcStatus {
55 #[prost(int32, tag = "1")]
56 pub code: i32,
57 #[prost(string, tag = "2")]
58 pub message: String,
59}
60
61fn is_json_content_type(content_type: Option<&ContentType>) -> bool {
62 match content_type {
63 None => false,
64 Some(ct) => {
65 let mime: mime::Mime = ct.clone().into();
66 mime.subtype() == mime::JSON
67 }
68 }
69}
70
71fn content_type_to_string(content_type: Option<&TypedHeader<ContentType>>) -> String {
72 content_type
73 .map(|h| h.0.to_string())
74 .unwrap_or_else(|| "not specified".to_string())
75}
76
77#[derive(Clone)]
78pub struct OtlpState {
79 pub with_metric_engine: bool,
80 pub handler: OpenTelemetryProtocolHandlerRef,
81}
82
83#[axum_macros::debug_handler]
84#[tracing::instrument(skip_all, fields(protocol = "otlp", request_type = "metrics"))]
85pub async fn metrics(
86 State(state): State<OtlpState>,
87 Extension(mut query_ctx): Extension<QueryContext>,
88 http_opts: OtlpMetricOptions,
89 content_type: Option<TypedHeader<ContentType>>,
90 bytes: Bytes,
91) -> Result<OtlpMetricsResponse> {
92 if is_json_content_type(content_type.as_ref().map(|h| &h.0)) {
93 return error::UnsupportedJsonContentTypeSnafu {}.fail();
94 }
95
96 let db = query_ctx.get_db_string();
97 query_ctx.set_channel(Channel::Otlp);
98
99 let _timer = crate::metrics::METRIC_HTTP_OPENTELEMETRY_METRICS_ELAPSED
100 .with_label_values(&[db.as_str()])
101 .start_timer();
102 let request = ExportMetricsServiceRequest::decode(bytes).with_context(|_| {
103 error::DecodeOtlpRequestSnafu {
104 content_type: content_type_to_string(content_type.as_ref()),
105 }
106 })?;
107
108 let OtlpState {
109 with_metric_engine,
110 handler,
111 } = state;
112
113 query_ctx.set_protocol_ctx(ProtocolCtx::OtlpMetric(OtlpMetricCtx {
114 promote_all_resource_attrs: http_opts.promote_all_resource_attrs,
115 resource_attrs: http_opts.resource_attrs,
116 promote_scope_attrs: http_opts.promote_scope_attrs,
117 with_metric_engine,
118 is_legacy: false,
120 resource_info: false,
121 metric_type: MetricType::Init,
122 metric_translation_strategy: http_opts.metric_translation_strategy,
123 }));
124 let query_ctx = Arc::new(query_ctx);
125
126 match handler.metrics(request, query_ctx).await {
127 Ok(outcome) => {
128 if outcome.accepted_data_points == 0 && outcome.rejected_data_points > 0 {
129 Ok(OtlpMetricsResponse::Failure(outcome))
130 } else if outcome.rejected_data_points > 0 || outcome.error_message.is_some() {
131 Ok(OtlpMetricsResponse::PartialSuccess(outcome))
132 } else {
133 Ok(OtlpMetricsResponse::FullSuccess(outcome))
134 }
135 }
136 Err(error::Error::InvalidOtlpMetricInput { reason }) => {
137 Ok(OtlpMetricsResponse::Failure(MetricsIngestOutcome {
138 error_message: Some(reason),
139 ..Default::default()
140 }))
141 }
142 Err(error) => Err(error),
143 }
144}
145
146#[axum_macros::debug_handler]
147#[tracing::instrument(skip_all, fields(protocol = "otlp", request_type = "traces"))]
148pub async fn traces(
149 State(state): State<OtlpState>,
150 TraceTableName(table_name): TraceTableName,
151 pipeline_info: PipelineInfo,
152 Extension(mut query_ctx): Extension<QueryContext>,
153 content_type: Option<TypedHeader<ContentType>>,
154 bytes: Bytes,
155) -> Result<OtlpTraceResponse> {
156 if is_json_content_type(content_type.as_ref().map(|h| &h.0)) {
157 return error::UnsupportedJsonContentTypeSnafu {}.fail();
158 }
159
160 let db = query_ctx.get_db_string();
161 let table_name = table_name.unwrap_or_else(|| TRACE_TABLE_NAME.to_string());
162
163 query_ctx.set_channel(Channel::Otlp);
164 query_ctx.set_extension(TRACE_TABLE_NAME_SESSION_KEY, &table_name);
165
166 let query_ctx = Arc::new(query_ctx);
167 let _timer = crate::metrics::METRIC_HTTP_OPENTELEMETRY_TRACES_ELAPSED
168 .with_label_values(&[db.as_str()])
169 .start_timer();
170 let request = ExportTraceServiceRequest::decode(bytes).with_context(|_| {
171 error::DecodeOtlpRequestSnafu {
172 content_type: content_type_to_string(content_type.as_ref()),
173 }
174 })?;
175
176 let pipeline = PipelineWay::from_name_and_default(
177 pipeline_info.pipeline_name.as_deref(),
178 pipeline_info.pipeline_version.as_deref(),
179 None,
180 )
181 .context(PipelineSnafu)?;
182
183 let pipeline_params = pipeline_info.pipeline_params;
184
185 let OtlpState { handler, .. } = state;
186
187 let pipeline_handler: Arc<dyn PipelineHandler + Send + Sync> = handler.clone();
190
191 handler
192 .traces(
193 pipeline_handler,
194 request,
195 pipeline,
196 pipeline_params,
197 table_name,
198 query_ctx,
199 )
200 .await
201 .map(|outcome| {
202 if outcome.accepted_spans == 0 && outcome.rejected_spans > 0 {
203 OtlpTraceResponse::Failure(outcome)
204 } else if outcome.rejected_spans > 0 || outcome.error_message.is_some() {
205 OtlpTraceResponse::PartialSuccess(outcome)
206 } else {
207 OtlpTraceResponse::FullSuccess(outcome)
208 }
209 })
210}
211
212#[axum_macros::debug_handler]
213#[tracing::instrument(skip_all, fields(protocol = "otlp", request_type = "logs"))]
214pub async fn logs(
215 State(state): State<OtlpState>,
216 Extension(mut query_ctx): Extension<QueryContext>,
217 pipeline_info: PipelineInfo,
218 LogTableName(tablename): LogTableName,
219 SelectInfoWrapper(select_info): SelectInfoWrapper,
220 content_type: Option<TypedHeader<ContentType>>,
221 bytes: Bytes,
222) -> Result<OtlpResponse<ExportLogsServiceResponse>> {
223 if is_json_content_type(content_type.as_ref().map(|h| &h.0)) {
224 return error::UnsupportedJsonContentTypeSnafu {}.fail();
225 }
226
227 let tablename = tablename.unwrap_or_else(|| "opentelemetry_logs".to_string());
228 let db = query_ctx.get_db_string();
229 query_ctx.set_channel(Channel::Otlp);
230 let query_ctx = Arc::new(query_ctx);
231 let _timer = METRIC_HTTP_OPENTELEMETRY_LOGS_ELAPSED
232 .with_label_values(&[db.as_str()])
233 .start_timer();
234 let request = ExportLogsServiceRequest::decode(bytes).with_context(|_| {
235 error::DecodeOtlpRequestSnafu {
236 content_type: content_type_to_string(content_type.as_ref()),
237 }
238 })?;
239
240 let pipeline = PipelineWay::from_name_and_default(
241 pipeline_info.pipeline_name.as_deref(),
242 pipeline_info.pipeline_version.as_deref(),
243 Some(PipelineWay::OtlpLogDirect(Box::new(select_info))),
244 )
245 .context(PipelineSnafu)?;
246 let pipeline_params = pipeline_info.pipeline_params;
247
248 let OtlpState { handler, .. } = state;
249
250 let pipeline_handler: Arc<dyn PipelineHandler + Send + Sync> = handler.clone();
253 handler
254 .logs(
255 pipeline_handler,
256 request,
257 pipeline,
258 pipeline_params,
259 tablename,
260 query_ctx,
261 )
262 .await
263 .map(|o| OtlpResponse {
264 resp_body: ExportLogsServiceResponse {
265 partial_success: None,
266 },
267 write_cost: o.iter().map(|o| o.meta.cost).sum(),
268 })
269}
270
271pub struct OtlpResponse<T: Message> {
272 resp_body: T,
273 write_cost: usize,
274}
275
276pub enum OtlpMetricsResponse {
277 FullSuccess(MetricsIngestOutcome),
278 PartialSuccess(MetricsIngestOutcome),
279 Failure(MetricsIngestOutcome),
280}
281
282impl IntoResponse for OtlpMetricsResponse {
283 fn into_response(self) -> axum::response::Response {
284 match self {
285 OtlpMetricsResponse::FullSuccess(outcome) => {
286 let mut header_map = write_cost_header_map(outcome.write_cost);
287 header_map.insert(header::CONTENT_TYPE, CONTENT_TYPE_PROTOBUF.clone());
288 let body = ExportMetricsServiceResponse {
289 partial_success: None,
290 };
291 (header_map, body.encode_to_vec()).into_response()
292 }
293 OtlpMetricsResponse::PartialSuccess(outcome) => {
294 let mut header_map = write_cost_header_map(outcome.write_cost);
295 header_map.insert(header::CONTENT_TYPE, CONTENT_TYPE_PROTOBUF.clone());
296 let body = ExportMetricsServiceResponse {
297 partial_success: Some(ExportMetricsPartialSuccess {
298 rejected_data_points: outcome.rejected_data_points,
299 error_message: outcome.error_message.unwrap_or_default(),
300 }),
301 };
302 (header_map, body.encode_to_vec()).into_response()
303 }
304 OtlpMetricsResponse::Failure(outcome) => {
305 let status = GoogleRpcStatus {
306 code: tonic::Code::InvalidArgument as i32,
307 message: outcome.error_message.unwrap_or_default(),
308 };
309 (
310 StatusCode::BAD_REQUEST,
311 [(header::CONTENT_TYPE, CONTENT_TYPE_PROTOBUF.as_ref())],
312 status.encode_to_vec(),
313 )
314 .into_response()
315 }
316 }
317 }
318}
319
320impl<T: Message> IntoResponse for OtlpResponse<T> {
321 fn into_response(self) -> axum::response::Response {
322 let mut header_map = write_cost_header_map(self.write_cost);
323 header_map.insert(header::CONTENT_TYPE, CONTENT_TYPE_PROTOBUF.clone());
324
325 (header_map, self.resp_body.encode_to_vec()).into_response()
326 }
327}
328
329pub enum OtlpTraceResponse {
330 FullSuccess(TraceIngestOutcome),
331 PartialSuccess(TraceIngestOutcome),
332 Failure(TraceIngestOutcome),
333}
334
335impl IntoResponse for OtlpTraceResponse {
336 fn into_response(self) -> axum::response::Response {
337 match self {
338 OtlpTraceResponse::FullSuccess(outcome) => {
339 let mut header_map = write_cost_header_map(outcome.write_cost);
340 header_map.insert(header::CONTENT_TYPE, CONTENT_TYPE_PROTOBUF.clone());
341 let body = ExportTraceServiceResponse {
342 partial_success: None,
343 };
344 (header_map, body.encode_to_vec()).into_response()
345 }
346 OtlpTraceResponse::PartialSuccess(outcome) => {
347 let mut header_map = write_cost_header_map(outcome.write_cost);
348 header_map.insert(header::CONTENT_TYPE, CONTENT_TYPE_PROTOBUF.clone());
349 let body = ExportTraceServiceResponse {
350 partial_success: outcome.error_message.map(|error_message| {
351 ExportTracePartialSuccess {
352 rejected_spans: outcome.rejected_spans as i64,
353 error_message,
354 }
355 }),
356 };
357 (header_map, body.encode_to_vec()).into_response()
358 }
359 OtlpTraceResponse::Failure(outcome) => {
360 let status = GoogleRpcStatus {
361 code: 0,
362 message: outcome.error_message.unwrap_or_default(),
363 };
364 (
365 StatusCode::BAD_REQUEST,
366 [(header::CONTENT_TYPE, CONTENT_TYPE_PROTOBUF.as_ref())],
367 status.encode_to_vec(),
368 )
369 .into_response()
370 }
371 }
372 }
373}
374
375#[cfg(test)]
376mod tests;