Skip to main content

frontend/instance/
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
15mod trace_ingest;
16pub mod trace_semconv;
17pub mod trace_types;
18
19use std::sync::Arc;
20
21use async_trait::async_trait;
22use auth::{
23    OTLP_WRITE, PermissionChecker, PermissionCheckerRef, PermissionReq, PermissionTableTarget,
24    PermissionTableTargets,
25};
26use client::Output;
27use common_catalog::consts::{trace_operations_table_name, trace_services_table_name};
28use common_error::ext::BoxedError;
29use common_query::prelude::GREPTIME_PHYSICAL_TABLE;
30use common_telemetry::tracing;
31use opentelemetry_proto::tonic::collector::logs::v1::ExportLogsServiceRequest;
32use opentelemetry_proto::tonic::collector::trace::v1::ExportTraceServiceRequest;
33use otel_arrow_rust::proto::opentelemetry::collector::metrics::v1::ExportMetricsServiceRequest;
34use pipeline::{GreptimePipelineParams, PipelineWay};
35use servers::error::{self, AuthSnafu, Result as ServerResult};
36use servers::http::prom_store::PHYSICAL_TABLE_PARAM;
37use servers::interceptor::{OpenTelemetryProtocolInterceptor, OpenTelemetryProtocolInterceptorRef};
38use servers::otlp;
39use servers::otlp::trace::span::TraceSpanGroup;
40use servers::query_handler::{
41    OpenTelemetryProtocolHandler, PipelineHandlerRef, TraceIngestOutcome,
42};
43use session::context::QueryContextRef;
44use snafu::ResultExt;
45use table::requests::{
46    OTLP_METRIC_COMPAT_KEY, OTLP_METRIC_COMPAT_PROM, SEMANTIC_PER_TABLE_INDEX_KEY,
47    SEMANTIC_SIGNAL_TYPE, SEMANTIC_SOURCE, SIGNAL_TYPE_LOG, SIGNAL_TYPE_METRIC,
48    SOURCE_OPENTELEMETRY,
49};
50
51use self::trace_ingest::trace_conventions;
52use crate::instance::Instance;
53use crate::metrics::{OTLP_LOGS_ROWS, OTLP_METRICS_ROWS};
54
55fn trace_permission_targets(
56    table_name: &str,
57    groups: &[TraceSpanGroup],
58    ctx: &QueryContextRef,
59) -> PermissionTableTargets {
60    let catalog = ctx.current_catalog();
61    let schema = ctx.current_schema();
62    if catalog.is_empty() || schema.is_empty() || table_name.is_empty() {
63        return PermissionTableTargets::Unresolved;
64    }
65
66    let has_spans = groups.iter().any(|group| !group.spans.is_empty());
67    if !has_spans {
68        return PermissionTableTargets::resolved(Vec::new());
69    }
70
71    let mut targets = vec![PermissionTableTarget::new(catalog, &schema, table_name)];
72    if groups
73        .iter()
74        .flat_map(|group| &group.spans)
75        .any(|span| span.service_name.is_some())
76    {
77        targets.extend([
78            PermissionTableTarget::new(catalog, &schema, trace_services_table_name(table_name)),
79            PermissionTableTarget::new(catalog, &schema, trace_operations_table_name(table_name)),
80        ]);
81    }
82
83    PermissionTableTargets::resolved(targets)
84}
85
86#[async_trait]
87impl OpenTelemetryProtocolHandler for Instance {
88    #[tracing::instrument(skip_all)]
89    async fn metrics(
90        &self,
91        request: ExportMetricsServiceRequest,
92        ctx: QueryContextRef,
93    ) -> ServerResult<Output> {
94        self.plugins
95            .get::<PermissionCheckerRef>()
96            .as_ref()
97            .check_permission(ctx.current_user(), PermissionReq::Action(OTLP_WRITE))
98            .context(AuthSnafu)?;
99
100        let interceptor_ref = self
101            .plugins
102            .get::<OpenTelemetryProtocolInterceptorRef<servers::error::Error>>();
103        interceptor_ref.pre_execute(ctx.clone())?;
104        let ctx = Arc::new(ctx.fork());
105
106        let input_names = request
107            .resource_metrics
108            .iter()
109            .flat_map(|r| r.scope_metrics.iter())
110            .flat_map(|s| s.metrics.iter().map(|m| m.name.clone()))
111            .collect::<Vec<_>>();
112
113        // See [`OtlpMetricCtx`] for details
114        let is_legacy = self.check_otlp_legacy(&input_names, &ctx).await?;
115
116        let mut metric_ctx = ctx
117            .protocol_ctx()
118            .get_otlp_metric_ctx()
119            .cloned()
120            .unwrap_or_default();
121        metric_ctx.is_legacy = is_legacy;
122
123        let (requests, rows, semantic_index) =
124            otlp::metrics::to_grpc_insert_requests(request, &mut metric_ctx)?;
125        self.check_row_insert_permission(&requests, &ctx, PermissionReq::Action(OTLP_WRITE))
126            .context(AuthSnafu)?;
127        self.cache_otlp_legacy(&input_names, &ctx, is_legacy)?;
128        OTLP_METRICS_ROWS.inc_by(rows as u64);
129
130        let ctx = {
131            let mut c = (*ctx).clone();
132            c.set_extension(SEMANTIC_SIGNAL_TYPE, SIGNAL_TYPE_METRIC);
133            c.set_extension(SEMANTIC_SOURCE, SOURCE_OPENTELEMETRY);
134            // Per-table metric specifics + resource/scope lineage ride this
135            // internal channel; the auto-create path folds them per schema and
136            // table name.
137            if let Some(index) = semantic_index.encode(&c.current_schema()) {
138                c.set_extension(SEMANTIC_PER_TABLE_INDEX_KEY, index);
139            }
140            if !is_legacy {
141                c.set_extension(OTLP_METRIC_COMPAT_KEY, OTLP_METRIC_COMPAT_PROM.to_string());
142            }
143            Arc::new(c)
144        };
145
146        // If the user uses the legacy path, it is by default without metric engine.
147        if metric_ctx.is_legacy || !metric_ctx.with_metric_engine {
148            self.handle_row_inserts(requests, ctx, false, false)
149                .await
150                .map_err(BoxedError::new)
151                .context(error::ExecuteGrpcQuerySnafu)
152        } else {
153            let physical_table = ctx
154                .extension(PHYSICAL_TABLE_PARAM)
155                .unwrap_or(GREPTIME_PHYSICAL_TABLE)
156                .to_string();
157            self.handle_metric_row_inserts(requests, ctx, physical_table.clone())
158                .await
159                .map_err(BoxedError::new)
160                .context(error::ExecuteGrpcQuerySnafu)
161        }
162    }
163
164    #[tracing::instrument(skip_all)]
165    async fn traces(
166        &self,
167        pipeline_handler: PipelineHandlerRef,
168        request: ExportTraceServiceRequest,
169        pipeline: PipelineWay,
170        pipeline_params: GreptimePipelineParams,
171        table_name: String,
172        ctx: QueryContextRef,
173    ) -> ServerResult<TraceIngestOutcome> {
174        self.plugins
175            .get::<PermissionCheckerRef>()
176            .as_ref()
177            .check_permission(ctx.current_user(), PermissionReq::Action(OTLP_WRITE))
178            .context(AuthSnafu)?;
179
180        let interceptor_ref = self
181            .plugins
182            .get::<OpenTelemetryProtocolInterceptorRef<servers::error::Error>>();
183        interceptor_ref.pre_execute(ctx.clone())?;
184        let ctx = Arc::new(ctx.fork());
185
186        // `schema_url` is consumed by `parse`, so derive conventions first.
187        let conventions = trace_conventions(&request);
188        let spans = otlp::trace::span::parse(request);
189        let targets = trace_permission_targets(&table_name, &spans, &ctx);
190        self.check_table_permission(&ctx, PermissionReq::Action(OTLP_WRITE), targets)
191            .context(AuthSnafu)?;
192        self.ingest_trace_spans(
193            pipeline_handler,
194            &pipeline,
195            &pipeline_params,
196            table_name,
197            spans,
198            &conventions,
199            ctx,
200        )
201        .await
202    }
203
204    #[tracing::instrument(skip_all)]
205    async fn logs(
206        &self,
207        pipeline_handler: PipelineHandlerRef,
208        request: ExportLogsServiceRequest,
209        pipeline: PipelineWay,
210        pipeline_params: GreptimePipelineParams,
211        table_name: String,
212        ctx: QueryContextRef,
213    ) -> ServerResult<Vec<Output>> {
214        self.plugins
215            .get::<PermissionCheckerRef>()
216            .as_ref()
217            .check_permission(ctx.current_user(), PermissionReq::Action(OTLP_WRITE))
218            .context(AuthSnafu)?;
219
220        let interceptor_ref = self
221            .plugins
222            .get::<OpenTelemetryProtocolInterceptorRef<servers::error::Error>>();
223        interceptor_ref.pre_execute(ctx.clone())?;
224        let ctx = Arc::new(ctx.fork());
225
226        // `as_req_iter` clones this ctx into each `temp_ctx`, so identity set here
227        // reaches the context that drives table auto-create.
228        let ctx = {
229            let mut c = (*ctx).clone();
230            c.set_extension(SEMANTIC_SIGNAL_TYPE, SIGNAL_TYPE_LOG);
231            c.set_extension(SEMANTIC_SOURCE, SOURCE_OPENTELEMETRY);
232            Arc::new(c)
233        };
234
235        let opt_req = otlp::logs::to_grpc_insert_requests(
236            request,
237            pipeline,
238            pipeline_params,
239            table_name,
240            &ctx,
241            pipeline_handler,
242        )
243        .await?;
244
245        let batches = opt_req.as_req_iter(ctx).collect::<Vec<_>>();
246        for (temp_ctx, requests) in &batches {
247            self.check_row_insert_permission(requests, temp_ctx, PermissionReq::Action(OTLP_WRITE))
248                .context(AuthSnafu)?;
249        }
250
251        let mut outputs = Vec::with_capacity(batches.len());
252        for (temp_ctx, requests) in batches {
253            let cnt = requests
254                .inserts
255                .iter()
256                .filter_map(|r| r.rows.as_ref().map(|r| r.rows.len()))
257                .sum::<usize>();
258
259            let o = self
260                .handle_log_inserts(requests, temp_ctx)
261                .await
262                .inspect(|_| OTLP_LOGS_ROWS.inc_by(cnt as u64))
263                .map_err(BoxedError::new)
264                .context(error::ExecuteGrpcQuerySnafu)?;
265            outputs.push(o);
266        }
267
268        Ok(outputs)
269    }
270}
271
272#[cfg(test)]
273mod tests {
274    use std::sync::Arc;
275
276    use auth::{PermissionTableTarget, PermissionTableTargets};
277    use opentelemetry_proto::tonic::collector::trace::v1::ExportTraceServiceRequest;
278    use opentelemetry_proto::tonic::common::v1::{AnyValue, KeyValue, any_value};
279    use opentelemetry_proto::tonic::resource::v1::Resource;
280    use opentelemetry_proto::tonic::trace::v1::{ResourceSpans, ScopeSpans, Span};
281    use session::context::QueryContext;
282
283    use super::trace_permission_targets;
284
285    #[test]
286    fn test_trace_permission_targets() {
287        let request = ExportTraceServiceRequest {
288            resource_spans: vec![ResourceSpans {
289                resource: Some(Resource {
290                    attributes: vec![KeyValue {
291                        key: "service.name".to_string(),
292                        value: Some(AnyValue {
293                            value: Some(any_value::Value::StringValue("frontend".to_string())),
294                        }),
295                        ..Default::default()
296                    }],
297                    ..Default::default()
298                }),
299                scope_spans: vec![ScopeSpans {
300                    spans: vec![Span::default()],
301                    ..Default::default()
302                }],
303                ..Default::default()
304            }],
305        };
306        let groups = servers::otlp::trace::span::parse(request);
307        let ctx = Arc::new(QueryContext::with("greptime", "public"));
308
309        assert_eq!(
310            PermissionTableTargets::Resolved(vec![
311                PermissionTableTarget::new("greptime", "public", "traces"),
312                PermissionTableTarget::new("greptime", "public", "traces_services"),
313                PermissionTableTarget::new("greptime", "public", "traces_operations"),
314            ]),
315            trace_permission_targets("traces", &groups, &ctx)
316        );
317        assert_eq!(
318            PermissionTableTargets::Unresolved,
319            trace_permission_targets("", &groups, &ctx)
320        );
321    }
322}