Skip to main content

servers/http/
otlp.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 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        // set by the frontend from its config
119        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    // here we use nightly feature `trait_upcasting` to convert handler to
188    // pipeline_handler
189    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    // here we use nightly feature `trait_upcasting` to convert handler to
251    // pipeline_handler
252    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;