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, warn};
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    MetricsIngestOutcome, 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, OTLP_RESOURCE_INFO_WRITE_ERRORS};
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<MetricsIngestOutcome> {
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        metric_ctx.resource_info = self.otlp_resource_info;
123
124        let otlp::metrics::MetricsConversion {
125            requests,
126            rows,
127            semantic_index,
128            resource_info,
129            mut outcome,
130        } = otlp::metrics::to_grpc_insert_requests(request, &mut metric_ctx)?;
131        if outcome.rejected_data_points > 0 {
132            warn!(
133                "Rejected {} OTLP metrics data points: {}",
134                outcome.rejected_data_points,
135                outcome.error_message.as_deref().unwrap_or_default()
136            );
137        }
138        if outcome.accepted_data_points == 0 {
139            return Ok(outcome);
140        }
141
142        self.check_row_insert_permission(&requests, &ctx, PermissionReq::Action(OTLP_WRITE))
143            .context(AuthSnafu)?;
144        self.cache_otlp_legacy(&input_names, &ctx, is_legacy)?;
145        OTLP_METRICS_ROWS.inc_by(rows as u64);
146
147        let ctx = {
148            let mut c = (*ctx).clone();
149            c.set_extension(SEMANTIC_SIGNAL_TYPE, SIGNAL_TYPE_METRIC);
150            c.set_extension(SEMANTIC_SOURCE, SOURCE_OPENTELEMETRY);
151            // Per-table metric specifics + resource/scope lineage ride this
152            // internal channel; the auto-create path folds them per schema and
153            // table name.
154            if let Some(index) = semantic_index.encode(&c.current_schema()) {
155                c.set_extension(SEMANTIC_PER_TABLE_INDEX_KEY, index);
156            }
157            if !is_legacy {
158                c.set_extension(OTLP_METRIC_COMPAT_KEY, OTLP_METRIC_COMPAT_PROM.to_string());
159            }
160            Arc::new(c)
161        };
162
163        // OTLP tables have one sample field in both the legacy and physical paths.
164        let output = if metric_ctx.is_legacy || !metric_ctx.with_metric_engine {
165            self.handle_row_inserts(requests, ctx.clone(), false, true)
166                .await
167                .map_err(BoxedError::new)
168                .context(error::ExecuteGrpcQuerySnafu)
169        } else {
170            let physical_table = ctx
171                .extension(PHYSICAL_TABLE_PARAM)
172                .unwrap_or(GREPTIME_PHYSICAL_TABLE)
173                .to_string();
174            self.handle_metric_row_inserts(requests, ctx.clone(), physical_table)
175                .await
176                .map_err(BoxedError::new)
177                .context(error::ExecuteGrpcQuerySnafu)
178        }?;
179        outcome.write_cost = output.meta.cost;
180
181        // Derived enrichment, written after the metric data is committed:
182        // failing here would make the client retry data the server already
183        // accepted, so every failure degrades to a warning instead.
184        if let Some(resource_info) = resource_info {
185            let written = match self.check_row_insert_permission(
186                &resource_info,
187                &ctx,
188                PermissionReq::Action(OTLP_WRITE),
189            ) {
190                Ok(_) => self
191                    .handle_row_inserts(resource_info, ctx, false, false)
192                    .await
193                    .map_err(BoxedError::new)
194                    .map_err(|e| e.to_string()),
195                Err(e) => Err(e.to_string()),
196            };
197            match written {
198                Ok(descriptor_output) => outcome.write_cost += descriptor_output.meta.cost,
199                Err(e) => {
200                    OTLP_RESOURCE_INFO_WRITE_ERRORS.inc();
201                    warn!("Failed to write the OTLP resource descriptor table: {e}");
202                    outcome.error_message.get_or_insert(format!(
203                        "metric data was accepted, but writing the resource \
204                         descriptor table `{}` failed: {e}",
205                        otlp::metrics::OTEL_RESOURCE_INFO_TABLE_NAME
206                    ));
207                }
208            }
209        }
210
211        Ok(outcome)
212    }
213
214    #[tracing::instrument(skip_all)]
215    async fn traces(
216        &self,
217        pipeline_handler: PipelineHandlerRef,
218        request: ExportTraceServiceRequest,
219        pipeline: PipelineWay,
220        pipeline_params: GreptimePipelineParams,
221        table_name: String,
222        ctx: QueryContextRef,
223    ) -> ServerResult<TraceIngestOutcome> {
224        self.plugins
225            .get::<PermissionCheckerRef>()
226            .as_ref()
227            .check_permission(ctx.current_user(), PermissionReq::Action(OTLP_WRITE))
228            .context(AuthSnafu)?;
229
230        let interceptor_ref = self
231            .plugins
232            .get::<OpenTelemetryProtocolInterceptorRef<servers::error::Error>>();
233        interceptor_ref.pre_execute(ctx.clone())?;
234        let ctx = Arc::new(ctx.fork());
235
236        // `schema_url` is consumed by `parse`, so derive conventions first.
237        let conventions = trace_conventions(&request);
238        let spans = otlp::trace::span::parse(request);
239        let targets = trace_permission_targets(&table_name, &spans, &ctx);
240        self.check_table_permission(&ctx, PermissionReq::Action(OTLP_WRITE), targets)
241            .context(AuthSnafu)?;
242        self.ingest_trace_spans(
243            pipeline_handler,
244            &pipeline,
245            &pipeline_params,
246            table_name,
247            spans,
248            &conventions,
249            ctx,
250        )
251        .await
252    }
253
254    #[tracing::instrument(skip_all)]
255    async fn logs(
256        &self,
257        pipeline_handler: PipelineHandlerRef,
258        request: ExportLogsServiceRequest,
259        pipeline: PipelineWay,
260        pipeline_params: GreptimePipelineParams,
261        table_name: String,
262        ctx: QueryContextRef,
263    ) -> ServerResult<Vec<Output>> {
264        self.plugins
265            .get::<PermissionCheckerRef>()
266            .as_ref()
267            .check_permission(ctx.current_user(), PermissionReq::Action(OTLP_WRITE))
268            .context(AuthSnafu)?;
269
270        let interceptor_ref = self
271            .plugins
272            .get::<OpenTelemetryProtocolInterceptorRef<servers::error::Error>>();
273        interceptor_ref.pre_execute(ctx.clone())?;
274        let ctx = Arc::new(ctx.fork());
275
276        // `as_req_iter` clones this ctx into each `temp_ctx`, so identity set here
277        // reaches the context that drives table auto-create.
278        let ctx = {
279            let mut c = (*ctx).clone();
280            c.set_extension(SEMANTIC_SIGNAL_TYPE, SIGNAL_TYPE_LOG);
281            c.set_extension(SEMANTIC_SOURCE, SOURCE_OPENTELEMETRY);
282            Arc::new(c)
283        };
284
285        let opt_req = otlp::logs::to_grpc_insert_requests(
286            request,
287            pipeline,
288            pipeline_params,
289            table_name,
290            &ctx,
291            pipeline_handler,
292        )
293        .await?;
294
295        let batches = opt_req.as_req_iter(ctx).collect::<Vec<_>>();
296        for (temp_ctx, requests) in &batches {
297            self.check_row_insert_permission(requests, temp_ctx, PermissionReq::Action(OTLP_WRITE))
298                .context(AuthSnafu)?;
299        }
300
301        let mut outputs = Vec::with_capacity(batches.len());
302        for (temp_ctx, requests) in batches {
303            let cnt = requests
304                .inserts
305                .iter()
306                .filter_map(|r| r.rows.as_ref().map(|r| r.rows.len()))
307                .sum::<usize>();
308
309            let o = self
310                .handle_log_inserts(requests, temp_ctx)
311                .await
312                .inspect(|_| OTLP_LOGS_ROWS.inc_by(cnt as u64))
313                .map_err(BoxedError::new)
314                .context(error::ExecuteGrpcQuerySnafu)?;
315            outputs.push(o);
316        }
317
318        Ok(outputs)
319    }
320}
321
322#[cfg(test)]
323mod tests {
324    use std::sync::Arc;
325
326    use auth::{PermissionTableTarget, PermissionTableTargets};
327    use opentelemetry_proto::tonic::collector::trace::v1::ExportTraceServiceRequest;
328    use opentelemetry_proto::tonic::common::v1::{AnyValue, KeyValue, any_value};
329    use opentelemetry_proto::tonic::resource::v1::Resource;
330    use opentelemetry_proto::tonic::trace::v1::{ResourceSpans, ScopeSpans, Span};
331    use session::context::QueryContext;
332
333    use super::trace_permission_targets;
334
335    #[test]
336    fn test_trace_permission_targets() {
337        let request = ExportTraceServiceRequest {
338            resource_spans: vec![ResourceSpans {
339                resource: Some(Resource {
340                    attributes: vec![KeyValue {
341                        key: "service.name".to_string(),
342                        value: Some(AnyValue {
343                            value: Some(any_value::Value::StringValue("frontend".to_string())),
344                        }),
345                        ..Default::default()
346                    }],
347                    ..Default::default()
348                }),
349                scope_spans: vec![ScopeSpans {
350                    spans: vec![Span::default()],
351                    ..Default::default()
352                }],
353                ..Default::default()
354            }],
355        };
356        let groups = servers::otlp::trace::span::parse(request);
357        let ctx = Arc::new(QueryContext::with("greptime", "public"));
358
359        assert_eq!(
360            PermissionTableTargets::Resolved(vec![
361                PermissionTableTarget::new("greptime", "public", "traces"),
362                PermissionTableTarget::new("greptime", "public", "traces_services"),
363                PermissionTableTarget::new("greptime", "public", "traces_operations"),
364            ]),
365            trace_permission_targets("traces", &groups, &ctx)
366        );
367        assert_eq!(
368            PermissionTableTargets::Unresolved,
369            trace_permission_targets("", &groups, &ctx)
370        );
371    }
372}