1use std::collections::{HashMap, HashSet};
16use std::sync::Arc;
17
18use async_trait::async_trait;
19use auth::{JAEGER_QUERY, PermissionReq, PermissionTableTarget, PermissionTableTargets};
20use catalog::CatalogManagerRef;
21use common_catalog::consts::{
22 TRACE_TABLE_NAME, trace_operations_table_name, trace_services_table_name,
23};
24use common_function::function::FunctionRef;
25use common_function::scalars::json::json_get::{
26 JsonGetBool, JsonGetFloat, JsonGetInt, JsonGetString, JsonGetWithType,
27};
28use common_function::scalars::udf::create_udf;
29use common_query::{Output, OutputData};
30use common_recordbatch::adapter::RecordBatchStreamAdapter;
31use common_recordbatch::util;
32use common_telemetry::warn;
33use datafusion::common::ScalarValue;
34use datafusion::dataframe::DataFrame;
35use datafusion::execution::SessionStateBuilder;
36use datafusion::execution::context::SessionContext;
37use datafusion::functions::core::expr_fn::coalesce;
38use datafusion::functions_window::expr_fn::row_number;
39use datafusion_expr::select_expr::SelectExpr;
40use datafusion_expr::{Expr, ExprFunctionExt, SortExpr, col, lit, lit_timestamp_nano, wildcard};
41use query::QueryEngineRef;
42use serde_json::Value as JsonValue;
43use servers::error::{
44 AuthSnafu, CollectRecordbatchSnafu, DataFusionSnafu, Result as ServerResult, TableNotFoundSnafu,
45};
46use servers::http::jaeger::{JAEGER_QUERY_TABLE_NAME_KEY, QueryTraceParams, TraceUserAgent};
47use servers::otlp::trace::{
48 DURATION_NANO_COLUMN, KEY_OTEL_STATUS_ERROR_KEY, RESOURCE_ATTRIBUTES_COLUMN,
49 SERVICE_NAME_COLUMN, SPAN_ATTRIBUTES_COLUMN, SPAN_KIND_COLUMN, SPAN_KIND_PREFIX,
50 SPAN_NAME_COLUMN, SPAN_STATUS_CODE, SPAN_STATUS_ERROR, TIMESTAMP_COLUMN, TRACE_ID_COLUMN,
51};
52use servers::query_handler::JaegerQueryHandler;
53use session::context::QueryContextRef;
54use snafu::{OptionExt, ResultExt};
55use table::TableRef;
56use table::table::adapter::DfTableProviderAdapter;
57
58use crate::instance::Instance;
59
60const DEFAULT_LIMIT: usize = 2000;
61const KEY_RN: &str = "greptime_rn";
62
63impl Instance {
64 async fn check_jaeger_query_permission(&self, ctx: &QueryContextRef) -> ServerResult<()> {
65 let table = ctx
66 .extension(JAEGER_QUERY_TABLE_NAME_KEY)
67 .unwrap_or(TRACE_TABLE_NAME);
68 let targets = PermissionTableTargets::resolved(vec![PermissionTableTarget::new(
69 ctx.current_catalog(),
70 ctx.current_schema(),
71 table,
72 )]);
73 let targets = self.resolve_query_permission_targets(targets, ctx).await?;
74 self.check_table_permission(ctx, PermissionReq::Action(JAEGER_QUERY), targets)
75 .context(AuthSnafu)?;
76 Ok(())
77 }
78}
79
80#[async_trait]
81impl JaegerQueryHandler for Instance {
82 async fn get_services(&self, ctx: QueryContextRef) -> ServerResult<Output> {
83 self.check_jaeger_query_permission(&ctx).await?;
84
85 Ok(query_trace_table(
87 ctx,
88 self,
89 vec![SelectExpr::from(col(SERVICE_NAME_COLUMN))],
90 vec![],
91 vec![],
92 None,
93 None,
94 vec![col(SERVICE_NAME_COLUMN)],
95 )
96 .await?)
97 }
98
99 async fn get_operations(
100 &self,
101 ctx: QueryContextRef,
102 service_name: &str,
103 span_kind: Option<&str>,
104 ) -> ServerResult<Output> {
105 self.check_jaeger_query_permission(&ctx).await?;
106
107 let mut filters = vec![col(SERVICE_NAME_COLUMN).eq(lit(service_name))];
108
109 if let Some(span_kind) = span_kind {
110 filters.push(col(SPAN_KIND_COLUMN).eq(lit(format!(
111 "{}{}",
112 SPAN_KIND_PREFIX,
113 span_kind.to_uppercase()
114 ))));
115 }
116
117 Ok(query_trace_table(
130 ctx,
131 self,
132 vec![
133 SelectExpr::from(col(SPAN_NAME_COLUMN)),
134 SelectExpr::from(col(SPAN_KIND_COLUMN)),
135 SelectExpr::from(col(SERVICE_NAME_COLUMN)),
136 SelectExpr::from(col(TIMESTAMP_COLUMN)),
137 ],
138 filters,
139 vec![col(SPAN_NAME_COLUMN).sort(true, false)], Some(DEFAULT_LIMIT),
141 None,
142 vec![col(SPAN_NAME_COLUMN), col(SPAN_KIND_COLUMN)],
143 )
144 .await?)
145 }
146
147 async fn get_trace(
148 &self,
149 ctx: QueryContextRef,
150 trace_id: &str,
151 start_time: Option<i64>,
152 end_time: Option<i64>,
153 limit: Option<usize>,
154 ) -> ServerResult<Output> {
155 self.check_jaeger_query_permission(&ctx).await?;
156
157 let selects = vec![wildcard()];
172
173 let mut filters = vec![col(TRACE_ID_COLUMN).eq(lit(trace_id))];
174
175 if let Some(start_time) = start_time {
176 filters.push(col(TIMESTAMP_COLUMN).gt_eq(lit_timestamp_nano(start_time)));
177 }
178
179 if let Some(end_time) = end_time {
180 filters.push(col(TIMESTAMP_COLUMN).lt_eq(lit_timestamp_nano(end_time)));
181 }
182
183 Ok(query_trace_table(
184 ctx,
185 self,
186 selects,
187 filters,
188 vec![col(TIMESTAMP_COLUMN).sort(false, false)], limit,
190 None,
191 vec![],
192 )
193 .await?)
194 }
195
196 async fn find_traces(
197 &self,
198 ctx: QueryContextRef,
199 query_params: QueryTraceParams,
200 ) -> ServerResult<Output> {
201 self.check_jaeger_query_permission(&ctx).await?;
202
203 let mut filters = vec![];
204
205 filters.push(col(SERVICE_NAME_COLUMN).eq(lit(query_params.service_name)));
207
208 if let Some(operation_name) = query_params.operation_name {
209 filters.push(col(SPAN_NAME_COLUMN).eq(lit(operation_name)));
210 }
211
212 if let Some(start_time) = query_params.start_time {
213 filters.push(col(TIMESTAMP_COLUMN).gt_eq(lit_timestamp_nano(start_time)));
214 }
215
216 if let Some(end_time) = query_params.end_time {
217 filters.push(col(TIMESTAMP_COLUMN).lt_eq(lit_timestamp_nano(end_time)));
218 }
219
220 if let Some(min_duration) = query_params.min_duration {
221 filters.push(col(DURATION_NANO_COLUMN).gt_eq(lit(min_duration)));
222 }
223
224 if let Some(max_duration) = query_params.max_duration {
225 filters.push(col(DURATION_NANO_COLUMN).lt_eq(lit(max_duration)));
226 }
227
228 let output = query_trace_table(
245 ctx.clone(),
246 self,
247 vec![wildcard()],
248 filters,
249 vec![],
250 Some(query_params.limit.unwrap_or(DEFAULT_LIMIT)),
251 query_params.tags,
252 vec![col(TRACE_ID_COLUMN)],
253 )
254 .await?;
255
256 let mut filters = vec![
269 col(TRACE_ID_COLUMN).in_list(
270 trace_ids_from_output(output)
271 .await?
272 .iter()
273 .map(lit)
274 .collect::<Vec<Expr>>(),
275 false,
276 ),
277 ];
278
279 if let Some(start_time) = query_params.start_time {
280 filters.push(col(TIMESTAMP_COLUMN).gt_eq(lit_timestamp_nano(start_time)));
281 }
282
283 if let Some(end_time) = query_params.end_time {
284 filters.push(col(TIMESTAMP_COLUMN).lt_eq(lit_timestamp_nano(end_time)));
285 }
286
287 match query_params.user_agent {
288 TraceUserAgent::Grafana => {
289 let table_name = ctx
293 .extension(JAEGER_QUERY_TABLE_NAME_KEY)
294 .unwrap_or(TRACE_TABLE_NAME);
295
296 let table = get_table(ctx.clone(), self.catalog_manager(), table_name).await?;
297
298 Ok(find_traces_rank_3(
299 table,
300 self.query_engine(),
301 filters,
302 vec![col(TIMESTAMP_COLUMN).sort(false, false)], )
304 .await?)
305 }
306 _ => {
307 Ok(query_trace_table(
309 ctx,
310 self,
311 vec![wildcard()],
312 filters,
313 vec![col(TIMESTAMP_COLUMN).sort(false, false)], None,
315 None,
316 vec![],
317 )
318 .await?)
319 }
320 }
321 }
322}
323
324#[allow(clippy::too_many_arguments)]
325async fn query_trace_table(
326 ctx: QueryContextRef,
327 instance: &Instance,
328 selects: Vec<SelectExpr>,
329 filters: Vec<Expr>,
330 sorts: Vec<SortExpr>,
331 limit: Option<usize>,
332 tags: Option<HashMap<String, JsonValue>>,
333 distincts: Vec<Expr>,
334) -> ServerResult<Output> {
335 let trace_table_name = ctx
336 .extension(JAEGER_QUERY_TABLE_NAME_KEY)
337 .unwrap_or(TRACE_TABLE_NAME);
338
339 let table_name = {
342 if match selects.as_slice() {
343 [SelectExpr::Expression(x)] => x == &col(SERVICE_NAME_COLUMN),
344 _ => false,
345 } {
346 &trace_services_table_name(trace_table_name)
347 } else if !distincts.is_empty()
348 && distincts.contains(&col(SPAN_NAME_COLUMN))
349 && distincts.contains(&col(SPAN_KIND_COLUMN))
350 {
351 &trace_operations_table_name(trace_table_name)
352 } else {
353 trace_table_name
354 }
355 };
356
357 let table = instance
358 .catalog_manager()
359 .table(
360 ctx.current_catalog(),
361 &ctx.current_schema(),
362 table_name,
363 Some(&ctx),
364 )
365 .await?
366 .with_context(|| TableNotFoundSnafu {
367 table: table_name,
368 catalog: ctx.current_catalog(),
369 schema: ctx.current_schema(),
370 })?;
371
372 let table_info = table.table_info();
373 let data_model = table_info.meta.options.data_model();
374
375 let col_names = table_info
377 .meta
378 .field_column_names()
379 .map(|s| format!("\"{}\"", s))
380 .collect::<HashSet<String>>();
381
382 let df_context = create_df_context(instance.query_engine())?;
383
384 let dataframe = df_context
385 .read_table(Arc::new(DfTableProviderAdapter::new(table)))
386 .context(DataFusionSnafu)?;
387
388 let dataframe = dataframe.select(selects).context(DataFusionSnafu)?;
389
390 let dataframe = filters
392 .into_iter()
393 .chain(tags.map_or(Ok(vec![]), |t| {
394 tags_filters(&dataframe, t, data_model, &col_names)
395 })?)
396 .try_fold(dataframe, |df, expr| {
397 df.filter(expr).context(DataFusionSnafu)
398 })?;
399
400 let dataframe = if !distincts.is_empty() {
402 dataframe
403 .distinct_on(distincts.clone(), distincts, None)
404 .context(DataFusionSnafu)?
405 } else {
406 dataframe
407 };
408
409 let dataframe = if !sorts.is_empty() {
411 dataframe.sort(sorts).context(DataFusionSnafu)?
412 } else {
413 dataframe
414 };
415
416 let dataframe = if let Some(limit) = limit {
418 dataframe.limit(0, Some(limit)).context(DataFusionSnafu)?
419 } else {
420 dataframe
421 };
422
423 let stream = dataframe.execute_stream().await.context(DataFusionSnafu)?;
425
426 let output = Output::new_with_stream(Box::pin(
427 RecordBatchStreamAdapter::try_new(stream).context(CollectRecordbatchSnafu)?,
428 ));
429
430 output
431 .map_dictionary_to_values()
432 .context(CollectRecordbatchSnafu)
433}
434
435async fn get_table(
436 ctx: QueryContextRef,
437 catalog_manager: &CatalogManagerRef,
438 table_name: &str,
439) -> ServerResult<TableRef> {
440 catalog_manager
441 .table(
442 ctx.current_catalog(),
443 &ctx.current_schema(),
444 table_name,
445 Some(&ctx),
446 )
447 .await?
448 .with_context(|| TableNotFoundSnafu {
449 table: table_name,
450 catalog: ctx.current_catalog(),
451 schema: ctx.current_schema(),
452 })
453}
454
455async fn find_traces_rank_3(
456 table: TableRef,
457 query_engine: &QueryEngineRef,
458 filters: Vec<Expr>,
459 sorts: Vec<SortExpr>,
460) -> ServerResult<Output> {
461 let df_context = create_df_context(query_engine)?;
462
463 let dataframe = df_context
464 .read_table(Arc::new(DfTableProviderAdapter::new(table)))
465 .context(DataFusionSnafu)?;
466
467 let dataframe = dataframe
468 .select(vec![wildcard()])
469 .context(DataFusionSnafu)?;
470
471 let dataframe = filters.into_iter().try_fold(dataframe, |df, expr| {
473 df.filter(expr).context(DataFusionSnafu)
474 })?;
475
476 let dataframe = if !sorts.is_empty() {
478 dataframe.sort(sorts).context(DataFusionSnafu)?
479 } else {
480 dataframe
481 };
482
483 let trace_id_col = vec![col(TRACE_ID_COLUMN)];
485 let timestamp_asc = vec![col(TIMESTAMP_COLUMN).sort(true, false)];
486
487 let dataframe = dataframe
488 .with_column(
489 KEY_RN,
490 row_number()
491 .partition_by(trace_id_col)
492 .order_by(timestamp_asc)
493 .build()
494 .context(DataFusionSnafu)?,
495 )
496 .context(DataFusionSnafu)?;
497
498 let dataframe = dataframe
499 .filter(col(KEY_RN).lt_eq(lit(3)))
500 .context(DataFusionSnafu)?;
501
502 let stream = dataframe.execute_stream().await.context(DataFusionSnafu)?;
504
505 let output = Output::new_with_stream(Box::pin(
506 RecordBatchStreamAdapter::try_new(stream).context(CollectRecordbatchSnafu)?,
507 ));
508
509 output
510 .map_dictionary_to_values()
511 .context(CollectRecordbatchSnafu)
512}
513
514fn create_df_context(query_engine: &QueryEngineRef) -> ServerResult<SessionContext> {
519 let df_context = SessionContext::new_with_state(
520 SessionStateBuilder::new_from_existing(query_engine.engine_state().session_state()).build(),
521 );
522
523 let udfs: Vec<FunctionRef> = vec![
525 Arc::new(JsonGetWithType::default()),
526 Arc::new(JsonGetInt::default()),
527 Arc::new(JsonGetFloat::default()),
528 Arc::new(JsonGetBool::default()),
529 Arc::new(JsonGetString::default()),
530 ];
531
532 for udf in udfs {
533 df_context.register_udf(create_udf(udf));
534 }
535
536 Ok(df_context)
537}
538
539fn json_tag_filters(
540 dataframe: &DataFrame,
541 tags: HashMap<String, JsonValue>,
542) -> ServerResult<Vec<Expr>> {
543 let mut filters = vec![];
544
545 for (key, value) in tags.iter() {
547 if let JsonValue::String(value) = value {
548 filters.push(
549 dataframe
550 .registry()
551 .udf(JsonGetString::NAME)
552 .context(DataFusionSnafu)?
553 .call(vec![
554 col(SPAN_ATTRIBUTES_COLUMN),
555 lit(format!("[\"{}\"]", key)),
556 ])
557 .eq(lit(value)),
558 );
559 }
560 if let JsonValue::Number(value) = value {
561 if value.is_i64() {
562 filters.push(
563 dataframe
564 .registry()
565 .udf(JsonGetInt::NAME)
566 .context(DataFusionSnafu)?
567 .call(vec![
568 col(SPAN_ATTRIBUTES_COLUMN),
569 lit(format!("[\"{}\"]", key)),
570 ])
571 .eq(lit(value.as_i64().unwrap())),
572 );
573 }
574 if value.is_f64() {
575 filters.push(
576 dataframe
577 .registry()
578 .udf(JsonGetFloat::NAME)
579 .context(DataFusionSnafu)?
580 .call(vec![
581 col(SPAN_ATTRIBUTES_COLUMN),
582 lit(format!("[\"{}\"]", key)),
583 ])
584 .eq(lit(value.as_f64().unwrap())),
585 );
586 }
587 }
588 if let JsonValue::Bool(value) = value {
589 filters.push(
590 dataframe
591 .registry()
592 .udf(JsonGetBool::NAME)
593 .context(DataFusionSnafu)?
594 .call(vec![
595 col(SPAN_ATTRIBUTES_COLUMN),
596 lit(format!("[\"{}\"]", key)),
597 ])
598 .eq(lit(*value)),
599 );
600 }
601 }
602
603 Ok(filters)
604}
605
606fn json2_tag_filters(
608 dataframe: &DataFrame,
609 tags: HashMap<String, JsonValue>,
610) -> ServerResult<Vec<Expr>> {
611 let get = dataframe
612 .registry()
613 .udf(JsonGetWithType::NAME)
614 .context(DataFusionSnafu)?;
615 tags.into_iter()
616 .map(|(key, value)| {
617 if key == KEY_OTEL_STATUS_ERROR_KEY && value == JsonValue::Bool(true) {
618 return Ok(col(SPAN_STATUS_CODE).eq(lit(SPAN_STATUS_ERROR)));
619 }
620 let path = lit(format!("$[{}]", JsonValue::String(key)));
622 let (value, value_type) = match value {
623 JsonValue::String(value) => (
624 ScalarValue::Utf8View(Some(value)),
625 ScalarValue::Utf8View(None),
626 ),
627 JsonValue::Bool(value) => (
628 ScalarValue::Boolean(Some(value)),
629 ScalarValue::Boolean(None),
630 ),
631 JsonValue::Number(value) => {
632 if let Some(value) = value.as_i64() {
633 (ScalarValue::Int64(Some(value)), ScalarValue::Int64(None))
634 } else if let Some(value) = value.as_u64() {
635 (ScalarValue::UInt64(Some(value)), ScalarValue::UInt64(None))
636 } else if let Some(value) = value.as_f64() {
637 (
638 ScalarValue::Float64(Some(value)),
639 ScalarValue::Float64(None),
640 )
641 } else {
642 return Ok(lit(false));
643 }
644 }
645 JsonValue::Null => (ScalarValue::Utf8View(None), ScalarValue::Utf8View(None)),
646 JsonValue::Array(_) | JsonValue::Object(_) => return Ok(lit(false)),
647 };
648 let attribute = coalesce(vec![
649 get.call(vec![
650 col(SPAN_ATTRIBUTES_COLUMN),
651 path.clone(),
652 lit(value_type.clone()),
653 ]),
654 get.call(vec![col(RESOURCE_ATTRIBUTES_COLUMN), path, lit(value_type)]),
655 ]);
656 Ok(if value.is_null() {
657 attribute.is_null()
658 } else {
659 attribute.eq(lit(value))
660 })
661 })
662 .collect()
663}
664
665#[inline]
668fn check_col_and_build_expr<F>(
669 span_key: String,
670 resource_key: String,
671 key: &str,
672 col_names: &HashSet<String>,
673 expr_builder: F,
674) -> Option<Expr>
675where
676 F: FnOnce(String) -> Expr,
677{
678 if col_names.contains(&span_key) {
679 return Some(expr_builder(span_key));
680 }
681 if col_names.contains(&resource_key) {
682 return Some(expr_builder(resource_key));
683 }
684 warn!("tag key {} not found in table columns", key);
685 None
686}
687
688fn flatten_tag_filters(
689 tags: HashMap<String, JsonValue>,
690 col_names: &HashSet<String>,
691) -> ServerResult<Vec<Expr>> {
692 let filters = tags
693 .into_iter()
694 .filter_map(|(key, value)| {
695 if key == KEY_OTEL_STATUS_ERROR_KEY && value == JsonValue::Bool(true) {
696 return Some(col(SPAN_STATUS_CODE).eq(lit(SPAN_STATUS_ERROR)));
697 }
698
699 let span_key = format!("\"span_attributes.{}\"", key);
701 let resource_key = format!("\"resource_attributes.{}\"", key);
702 match value {
703 JsonValue::String(value) => {
704 check_col_and_build_expr(span_key, resource_key, &key, col_names, |k| {
705 col(k).eq(lit(value))
706 })
707 }
708 JsonValue::Number(value) => {
709 if value.is_f64() {
710 let value = value.as_f64().unwrap();
712 check_col_and_build_expr(span_key, resource_key, &key, col_names, |k| {
713 col(k).eq(lit(value))
714 })
715 } else {
716 let value = value.as_i64().unwrap();
717 check_col_and_build_expr(span_key, resource_key, &key, col_names, |k| {
718 col(k).eq(lit(value))
719 })
720 }
721 }
722 JsonValue::Bool(value) => {
723 check_col_and_build_expr(span_key, resource_key, &key, col_names, |k| {
724 col(k).eq(lit(value))
725 })
726 }
727 JsonValue::Null => {
728 check_col_and_build_expr(span_key, resource_key, &key, col_names, |k| {
729 col(k).is_null()
730 })
731 }
732 JsonValue::Array(_value) => None,
734 JsonValue::Object(_value) => None,
735 }
736 })
737 .collect();
738 Ok(filters)
739}
740
741fn tags_filters(
742 dataframe: &DataFrame,
743 tags: HashMap<String, JsonValue>,
744 data_model: Option<&str>,
745 col_names: &HashSet<String>,
746) -> ServerResult<Vec<Expr>> {
747 match data_model {
748 Some(table::requests::TABLE_DATA_MODEL_TRACE_V1) => flatten_tag_filters(tags, col_names),
749 Some(table::requests::TABLE_DATA_MODEL_TRACE_V2) => json2_tag_filters(dataframe, tags),
750 _ => json_tag_filters(dataframe, tags),
751 }
752}
753
754async fn trace_ids_from_output(output: Output) -> ServerResult<Vec<String>> {
756 if let OutputData::Stream(stream) = output.data {
757 let schema = stream.schema().clone();
758 let recordbatches = util::collect(stream)
759 .await
760 .context(CollectRecordbatchSnafu)?;
761
762 if !recordbatches.is_empty()
764 && schema.num_columns() == 1
765 && schema.contains_column(TRACE_ID_COLUMN)
766 {
767 let mut trace_ids = vec![];
768 for recordbatch in recordbatches {
769 recordbatch
770 .iter_column_as_string(0)
771 .flatten()
772 .for_each(|x| trace_ids.push(x));
773 }
774
775 return Ok(trace_ids);
776 }
777 }
778
779 Ok(vec![])
780}