1mod 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 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 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 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 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 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}