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, 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 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 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 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 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 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 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}